From e8d98abccbf7bcf9e1267f1fe0c6550e6eb5f926 Mon Sep 17 00:00:00 2001 From: Maurus Decimus <11444311+mdecimus@users.noreply.github.com> Date: Thu, 27 Aug 2026 19:15:03 +0200 Subject: [PATCH] MS Exchange Graph: Import the default Contacts folder, recover series exceptions (fixes #39) --- CHANGELOG.md | 10 + Cargo.lock | 46 +- Cargo.toml | 2 +- src/cli.rs | 19 +- src/db/exchange_graph_ids.rs | 1 + src/db/init.rs | 36 ++ src/db/schema.sql | 3 +- src/exchange_graph/api.rs | 68 ++- src/exchange_graph/calendar_map.rs | 19 +- src/exchange_graph/client.rs | 11 +- src/exchange_graph/recurrence.rs | 33 +- src/exchange_graph/types.rs | 22 +- src/sync/import_exchange_graph.rs | 1 + src/sync/import_exchange_graph/calendar.rs | 141 +++++- src/sync/import_exchange_graph/contacts.rs | 12 +- src/sync/import_exchange_graph/coordinator.rs | 33 +- src/sync/import_exchange_graph/files.rs | 438 +++++++++++++++++ src/sync/import_exchange_graph/folders.rs | 68 ++- src/sync/import_exchange_graph/messages.rs | 66 ++- tests/mock_exchange_graph.rs | 461 ++++++++++++++++-- 20 files changed, 1370 insertions(+), 120 deletions(-) create mode 100644 src/sync/import_exchange_graph/files.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index de8b810..520c2fd 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,16 @@ All notable changes to this project will be documented in this file. This project adheres to [Semantic Versioning](http://semver.org/). +## [1.0.10] - 2026-08-27 + +### Added +- MS Exchange Graph: Import the default Contacts folder, recover series exceptions (fixes #39) +- OneDrive files import + +### Changed + +### Fixed + ## [1.0.9] - 2026-08-22 ### Added diff --git a/Cargo.lock b/Cargo.lock index f7c3158..69ea5c0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -157,7 +157,7 @@ checksum = "82f6aeea286b8eb4dd3431a1be1b59d290ace00f5bfd8e2a159bc2a05e2c1667" dependencies = [ "proc-macro2", "quote", - "syn 3.0.3", + "syn 3.0.4", ] [[package]] @@ -371,9 +371,9 @@ checksum = "fc652a48c352aef3ea3aed32080501cf3ef6ed5da78602a020c991775b0aff04" [[package]] name = "calcard" -version = "0.3.11" +version = "0.3.13" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f908fcb612ff8729e4302e071562bfbb0f7cdeb1b791e06024e7f11eedb83d9c" +checksum = "75b779382e675380a1ff8a4873acee5ba60158be7a00458fb98e6b85aad1f2ee" dependencies = [ "ahash", "chrono", @@ -471,7 +471,7 @@ dependencies = [ "heck", "proc-macro2", "quote", - "syn 3.0.3", + "syn 3.0.4", ] [[package]] @@ -506,9 +506,9 @@ dependencies = [ [[package]] name = "combine" -version = "4.6.7" +version = "4.6.8" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ba5a308b75df32fe02788e748662718f03fde005016435c444eea572398219fd" +checksum = "cfc320937d09e6de266b31b9afb480f197d7a861be86be7cb2ea7e5d1bfffc5e" dependencies = [ "bytes", "memchr", @@ -576,9 +576,9 @@ dependencies = [ [[package]] name = "crc32fast" -version = "1.5.0" +version = "1.5.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9481c1c90cbf2ac953f07c8d4a58aa3945c425b7185c9154d67a65e4230da511" +checksum = "8498c871161e1742aaa9d52551b2d6ebdd4c3d45a3be423e3728f33b955be550" dependencies = [ "cfg-if", ] @@ -680,7 +680,7 @@ checksum = "c6232dd377dcc64799954cbd3a9bb882e9cdc1308ccd87b1c098f1fb2eaf82a8" dependencies = [ "proc-macro2", "quote", - "syn 3.0.3", + "syn 3.0.4", ] [[package]] @@ -875,7 +875,7 @@ checksum = "9fb9654ba8355388abeb8dcb4fc62f511300867002afc858860463bdd9fe0c44" dependencies = [ "proc-macro2", "quote", - "syn 3.0.3", + "syn 3.0.4", ] [[package]] @@ -944,9 +944,9 @@ dependencies = [ [[package]] name = "h2" -version = "0.4.18" +version = "0.4.19" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "839c0e8a181239723652be9062bb56ca5bf5f64011f73b623f6f4fc59086a228" +checksum = "ef8e5e5a340588f4452631496976cf8636d4a7ecf600239fdc27615d2530bc16" dependencies = [ "atomic-waker", "bytes", @@ -2025,7 +2025,7 @@ checksum = "92ecd8964f8453721699a1ed72037b0db49ce2f5a5138486ee89bed6f67cdf3a" dependencies = [ "proc-macro2", "quote", - "syn 3.0.3", + "syn 3.0.4", ] [[package]] @@ -2316,7 +2316,7 @@ checksum = "e7a5d71263a5a7d47b41f6b3f06ba276f10cc18b0931f1799f710578e2309348" dependencies = [ "proc-macro2", "quote", - "syn 3.0.3", + "syn 3.0.4", ] [[package]] @@ -2340,7 +2340,7 @@ checksum = "8d3b1629de253c70a0508c3899572da79ca359fdab27c7920ff00406df418906" dependencies = [ "proc-macro2", "quote", - "syn 3.0.3", + "syn 3.0.4", ] [[package]] @@ -2532,9 +2532,9 @@ dependencies = [ [[package]] name = "syn" -version = "3.0.3" +version = "3.0.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "53e9bae58849f64dfa4f5d5ae372c8341f7305f82a3868709269343628b659a3" +checksum = "e6275cddf4610d1775e6d1fe9469b2e77d0f39fd98fb7450901b821e0c53649f" dependencies = [ "proc-macro2", "quote", @@ -2619,7 +2619,7 @@ checksum = "bc04cd3e1236dd4a98afca4569f2deb3f120e5422a4023be2cb683f8486292af" dependencies = [ "proc-macro2", "quote", - "syn 3.0.3", + "syn 3.0.4", ] [[package]] @@ -2702,7 +2702,7 @@ checksum = "78773a2a397f451582ce068015985c33193cf6dea8b74d2a639fe457b2f07b0e" dependencies = [ "proc-macro2", "quote", - "syn 3.0.3", + "syn 3.0.4", ] [[package]] @@ -2925,16 +2925,16 @@ checksum = "06abde3611657adf66d383f00b093d7faecc7fa57071cce2578660c9f1010821" [[package]] name = "uuid" -version = "1.25.0" +version = "1.26.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f053576934f05a761a402421fbbe3d425d9366f75f978806a037b3ca481abecc" +checksum = "b5772d71c9be8a8a6ac2117d949c5b224c1b72241bb611d9a3012edcf8af7812" dependencies = [ "sha1_smol", ] [[package]] name = "vandelay" -version = "1.0.9" +version = "1.0.10" dependencies = [ "base64 0.23.1", "blake3", @@ -3368,7 +3368,7 @@ checksum = "34df6fc39dbd26ddc9c10e6a2984476e13acce22e64e4487636ef494369225da" dependencies = [ "proc-macro2", "quote", - "syn 3.0.3", + "syn 3.0.4", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 5cc9143..80a4d9b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,7 +1,7 @@ [package] name = "vandelay" description = "JMAP account migration utility" -version = "1.0.9" +version = "1.0.10" authors = ["Stalwart Labs LLC "] license = "Apache-2.0 OR MIT" repository = "https://github.com/stalwartlabs/vandelay" diff --git a/src/cli.rs b/src/cli.rs index 422d4a5..383d903 100644 --- a/src/cli.rs +++ b/src/cli.rs @@ -1132,7 +1132,7 @@ pub struct ExchangeGraphImportArgs { #[arg( long, value_name = "LIST", - help = "Comma-separated surface list: mail | calendar | contacts (default: all three)" + help = "Comma-separated surface list: mail | calendar | contacts | files (default: all)" )] objects: Option, @@ -1161,6 +1161,19 @@ pub struct ExchangeGraphImportArgs { )] top: usize, + #[arg( + long, + value_name = "N", + default_value_t = 5, + help = "Years either side of today to scan for edited recurrence occurrences (max 5)", + long_help = "Years either side of today to scan for edited recurrence occurrences.\n \ + Graph reports a series exception only inside an expanded calendar view, and caps\n \ + any single view at five years, so vandelay scans [today - N, today] and\n \ + [today, today + N]. An occurrence edited outside that span imports with the\n \ + series default instead of its edit." + )] + exception_window_years: i32, + #[arg( long, value_name = "URL", @@ -1225,6 +1238,7 @@ fn resolve_exchange_graph_import(args: ExchangeGraphImportArgs) -> Result Result<(), OpenError> { let tx = conn.unchecked_transaction()?; tx.execute_batch(SCHEMA_SQL)?; ensure_calendar_events_data_type(&tx)?; + ensure_graph_ids_accept_file_nodes(&tx)?; tx.commit()?; Ok(()) } @@ -36,6 +37,41 @@ fn ensure_calendar_events_data_type(conn: &Connection) -> Result<(), OpenError> Ok(()) } +fn ensure_graph_ids_accept_file_nodes(conn: &Connection) -> Result<(), OpenError> { + let sql: Option = conn + .query_row( + "SELECT sql FROM sqlite_master WHERE type = 'table' \ + AND name = 'sync_id_exchange_graph'", + [], + |row| row.get(0), + ) + .ok(); + let Some(sql) = sql else { return Ok(()) }; + if sql.contains("'filenode'") { + return Ok(()); + } + conn.execute_batch( + "ALTER TABLE sync_id_exchange_graph RENAME TO sync_id_exchange_graph_old; + CREATE TABLE sync_id_exchange_graph ( + source_id INTEGER NOT NULL REFERENCES sources(id) ON DELETE CASCADE, + type_name TEXT NOT NULL CHECK (type_name IN ( + 'mailbox','email', + 'calendar','calendarevent', + 'addressbook','contactcard', + 'filenode')), + graph_id TEXT NOT NULL, + local_id INTEGER NOT NULL, + PRIMARY KEY (source_id, type_name, graph_id), + UNIQUE (source_id, type_name, local_id) + ); + INSERT INTO sync_id_exchange_graph SELECT * FROM sync_id_exchange_graph_old; + DROP TABLE sync_id_exchange_graph_old; + CREATE INDEX IF NOT EXISTS sync_id_exchange_graph_type_idx + ON sync_id_exchange_graph (source_id, type_name);", + )?; + Ok(()) +} + fn apply_pragmas(conn: &Connection) -> Result<(), OpenError> { conn.pragma_update(None, "journal_mode", "WAL")?; conn.pragma_update(None, "foreign_keys", "ON")?; diff --git a/src/db/schema.sql b/src/db/schema.sql index ca0c041..5ed1f37 100644 --- a/src/db/schema.sql +++ b/src/db/schema.sql @@ -118,7 +118,8 @@ CREATE TABLE IF NOT EXISTS sync_id_exchange_graph ( type_name TEXT NOT NULL CHECK (type_name IN ( 'mailbox','email', 'calendar','calendarevent', - 'addressbook','contactcard')), + 'addressbook','contactcard', + 'filenode')), graph_id TEXT NOT NULL, local_id INTEGER NOT NULL, PRIMARY KEY (source_id, type_name, graph_id), diff --git a/src/exchange_graph/api.rs b/src/exchange_graph/api.rs index a7d8d2f..9756493 100644 --- a/src/exchange_graph/api.rs +++ b/src/exchange_graph/api.rs @@ -16,6 +16,18 @@ pub const PREFER_TIMEZONE_UTC: &str = "outlook.timezone=\"UTC\""; pub const PREFER_BODY_TEXT: &str = "outlook.body-content-type=\"text\""; pub const PREFER_BODY_HTML: &str = "outlook.body-content-type=\"html\""; +pub const EXCEPTION_SELECT: &str = "id,type,seriesMasterId,originalStart,iCalUId,subject,body,\ + start,end,isAllDay,isCancelled,isDraft,sensitivity,importance,showAs,categories,locations,\ + organizer,attendees,createdDateTime,lastModifiedDateTime,isReminderOn,\ + reminderMinutesBeforeStart,originalStartTimeZone"; + +pub const EXCEPTION_WINDOW_MAX_DAYS: i64 = 1825; + +pub const DRIVE_ITEM_SELECT: &str = + "id,name,size,folder,file,package,remoteItem,createdDateTime,lastModifiedDateTime"; + +pub const MESSAGE_STATE_SELECT: &str = "id,isRead,isDraft,isReadReceiptRequested,flag,categories"; + #[derive(Debug, Clone)] pub struct Endpoints { pub api_base: String, @@ -68,7 +80,7 @@ impl Endpoints { pub fn folder_messages_ids(&self, folder_id: &str, top: usize) -> String { format!( - "{}/mailFolders/{}/messages?$top={top}&$select=id", + "{}/mailFolders/{}/messages?$top={top}&$select={MESSAGE_STATE_SELECT}", self.me_or_user(), url_escape(folder_id) ) @@ -106,10 +118,64 @@ impl Endpoints { format!("{}/events/{}", self.me_or_user(), url_escape(event_id)) } + pub fn calendar_exceptions( + &self, + calendar_id: &str, + window_start: &str, + window_end: &str, + top: usize, + ) -> String { + format!( + "{}/calendars/{}/calendarView?startDateTime={window_start}&endDateTime={window_end}\ + &$filter=type%20eq%20%27exception%27&$top={top}&$select={EXCEPTION_SELECT}", + self.me_or_user(), + url_escape(calendar_id) + ) + } + + pub fn drive_root(&self) -> String { + format!("{}/drive/root?$select=id,name", self.me_or_user()) + } + + pub fn drive_children(&self, item_id: &str, top: usize) -> String { + format!( + "{}/drive/items/{}/children?$top={top}&$select={DRIVE_ITEM_SELECT}", + self.me_or_user(), + url_escape(item_id) + ) + } + + pub fn drive_item_content(&self, item_id: &str) -> String { + format!( + "{}/drive/items/{}/content", + self.me_or_user(), + url_escape(item_id) + ) + } + pub fn contact_folders(&self, top: usize) -> String { format!("{}/contactFolders?$top={top}", self.me_or_user()) } + pub fn default_contact_folder(&self) -> String { + format!("{}/contactFolders/contacts", self.me_or_user()) + } + + pub fn contact_folder(&self, folder_id: &str) -> String { + format!( + "{}/contactFolders/{}", + self.me_or_user(), + url_escape(folder_id) + ) + } + + pub fn any_contact_parent(&self) -> String { + format!( + "{}/contacts?$top=1&$select=id,parentFolderId", + self.me_or_user() + ) + } + pub fn contact_folder_children(&self, folder_id: &str, top: usize) -> String { format!( "{}/contactFolders/{}/childFolders?$top={top}", diff --git a/src/exchange_graph/calendar_map.rs b/src/exchange_graph/calendar_map.rs index 690dc2f..c741020 100644 --- a/src/exchange_graph/calendar_map.rs +++ b/src/exchange_graph/calendar_map.rs @@ -320,9 +320,15 @@ pub fn convert_event( .and_then(Value::as_array) && !cancels.is_empty() { + let start_time = card + .get("start") + .and_then(Value::as_str) + .and_then(|s| s.split_once('T').map(|(_, time)| time.to_owned())); let mut overrides = Map::new(); for entry in cancels.iter().filter_map(Value::as_str) { - overrides.insert(entry.to_owned(), json!({"excluded": true})); + if let Some(key) = cancelled_occurrence_key(entry, start_time.as_deref()) { + overrides.insert(key, json!({"excluded": true})); + } } if !overrides.is_empty() { card.insert("recurrenceOverrides".to_owned(), Value::Object(overrides)); @@ -340,6 +346,17 @@ pub fn convert_event( }) } +fn cancelled_occurrence_key(entry: &str, start_time: Option<&str>) -> Option { + let date = entry.rsplit('.').next().filter(|d| d.len() == 10)?; + if !date.as_bytes().iter().enumerate().all(|(i, b)| match i { + 4 | 7 => *b == b'-', + _ => b.is_ascii_digit(), + }) { + return None; + } + Some(format!("{date}T{}", start_time.unwrap_or("00:00:00"))) +} + fn extract_local_datetime(slot: Option<&Value>) -> Option { let slot = slot?; let dt = slot.get("dateTime").and_then(Value::as_str)?; diff --git a/src/exchange_graph/client.rs b/src/exchange_graph/client.rs index 104f61c..7840b9e 100644 --- a/src/exchange_graph/client.rs +++ b/src/exchange_graph/client.rs @@ -28,6 +28,7 @@ const LONG_RETRY_THRESHOLD: Duration = Duration::from_secs(10); pub enum Accept { Json, Text, + Binary, } impl Accept { @@ -35,6 +36,7 @@ impl Accept { match self { Accept::Json => "application/json", Accept::Text => "text/plain", + Accept::Binary => "application/octet-stream", } } } @@ -411,10 +413,17 @@ struct RetryLog<'a> { } fn warn_on_redirect(logger: &Logger, method: &str, requested: &str, final_uri: &Uri) { + if requested.ends_with("/content") { + return; + } let landed = final_uri.to_string(); if cross_host(requested, &landed) { + let host = url::Url::parse(&landed) + .ok() + .and_then(|u| u.host_str().map(str::to_owned)) + .unwrap_or_else(|| "another host".to_owned()); logger.warn(&format!( - "{method} {requested} was redirected across hosts to {landed}; verify the Graph endpoint host is reachable directly without redirection" + "{method} {requested} was redirected across hosts to {host}; verify the Graph endpoint host is reachable directly without redirection" )); } } diff --git a/src/exchange_graph/recurrence.rs b/src/exchange_graph/recurrence.rs index 72692ea..43568eb 100644 --- a/src/exchange_graph/recurrence.rs +++ b/src/exchange_graph/recurrence.rs @@ -69,7 +69,9 @@ pub fn convert_patterned_recurrence_rule( ); } - if let Some(index) = pattern.get("index").and_then(Value::as_str) + let pattern_type = pattern.get("type").and_then(Value::as_str).unwrap_or(""); + if matches!(pattern_type, "relativeMonthly" | "relativeYearly") + && let Some(index) = pattern.get("index").and_then(Value::as_str) && let Some(setpos) = set_position_for(index) { rule.insert( @@ -205,6 +207,35 @@ mod tests { assert_eq!(rule["byDay"][0]["day"], "th"); } + #[test] + fn index_is_ignored_unless_the_pattern_is_relative() { + for kind in ["daily", "weekly", "absoluteMonthly", "absoluteYearly"] { + let pr = json!({ + "pattern": {"type": kind, "interval": 1, "index": "first", + "daysOfWeek": ["monday"], "dayOfMonth": 1, "month": 1}, + "range": {"type": "noEnd"} + }); + let out = convert_patterned_recurrence(&pr).unwrap(); + assert!( + out.get("bySetPosition").is_none(), + "{kind} carries Graph's default index=first, which must not become bySetPosition" + ); + } + } + + #[test] + fn index_survives_for_relative_patterns() { + for kind in ["relativeMonthly", "relativeYearly"] { + let pr = json!({ + "pattern": {"type": kind, "interval": 1, "index": "third", + "daysOfWeek": ["thursday"]}, + "range": {"type": "noEnd"} + }); + let out = convert_patterned_recurrence(&pr).unwrap(); + assert_eq!(out["bySetPosition"][0], 3, "{kind} keeps its index"); + } + } + #[test] fn absolute_yearly_with_month_and_day() { let pr = json!({ diff --git a/src/exchange_graph/types.rs b/src/exchange_graph/types.rs index b15e70d..31bcde0 100644 --- a/src/exchange_graph/types.rs +++ b/src/exchange_graph/types.rs @@ -64,6 +64,7 @@ pub struct Surfaces { pub mail: bool, pub calendar: bool, pub contacts: bool, + pub files: bool, } impl Default for Surfaces { @@ -77,12 +78,14 @@ impl Surfaces { mail: true, calendar: true, contacts: true, + files: true, }; pub const NONE: Surfaces = Surfaces { mail: false, calendar: false, contacts: false, + files: false, }; pub fn parse_list(list: &str) -> Result { @@ -96,9 +99,10 @@ impl Surfaces { "mail" => selected.mail = true, "calendar" => selected.calendar = true, "contacts" => selected.contacts = true, + "files" => selected.files = true, _ => { return Err(Error::Usage(format!( - "unknown surface: {token} (valid: mail, calendar, contacts)" + "unknown surface: {token} (valid: mail, calendar, contacts, files)" ))); } } @@ -158,7 +162,7 @@ mod tests { #[test] fn surfaces_default_is_all_three() { let all = Surfaces::default(); - assert!(all.mail && all.calendar && all.contacts); + assert!(all.mail && all.calendar && all.contacts && all.files); } #[test] @@ -168,7 +172,8 @@ mod tests { Surfaces { mail: true, calendar: false, - contacts: false + contacts: false, + files: false } ); assert_eq!( @@ -176,7 +181,8 @@ mod tests { Surfaces { mail: false, calendar: true, - contacts: false + contacts: false, + files: false } ); assert_eq!( @@ -184,7 +190,8 @@ mod tests { Surfaces { mail: false, calendar: false, - contacts: true + contacts: true, + files: false } ); } @@ -192,7 +199,7 @@ mod tests { #[test] fn surfaces_parse_is_case_insensitive() { assert_eq!( - Surfaces::parse_list("Mail,CALENDAR,Contacts").unwrap(), + Surfaces::parse_list("Mail,CALENDAR,Contacts,Files").unwrap(), Surfaces::ALL ); } @@ -204,7 +211,8 @@ mod tests { Surfaces { mail: true, calendar: false, - contacts: true + contacts: true, + files: false } ); } diff --git a/src/sync/import_exchange_graph.rs b/src/sync/import_exchange_graph.rs index 17c3c05..1d29802 100644 --- a/src/sync/import_exchange_graph.rs +++ b/src/sync/import_exchange_graph.rs @@ -7,6 +7,7 @@ pub mod calendar; pub mod contacts; pub mod coordinator; +pub mod files; pub mod folders; pub mod messages; diff --git a/src/sync/import_exchange_graph/calendar.rs b/src/sync/import_exchange_graph/calendar.rs index e20c830..70b578e 100644 --- a/src/sync/import_exchange_graph/calendar.rs +++ b/src/sync/import_exchange_graph/calendar.rs @@ -57,7 +57,9 @@ pub fn reconcile_all( _ => { if let Some(id) = stub.get("id").and_then(Value::as_str) { server_total.insert(id.to_owned()); - if !local.contains_key(id) && planned.insert(id.to_owned()) { + if local.contains_key(id) { + counts.fetched += 1; + } else if planned.insert(id.to_owned()) { want_ids.push(id.to_owned()); } } @@ -79,19 +81,23 @@ pub fn reconcile_all( let fetched = fetch_events(ctx, &want_ids); let mut masters: Vec<(String, ConvertedEvent)> = Vec::new(); let mut exceptions: Vec<(String, ConvertedEvent)> = Vec::new(); + let mut any_series = false; for (graph_id, result) in fetched { match result { - Ok(raw) => match convert_event(&raw, None) { - Ok(c) => match c.event_type { - EventType::Exception => exceptions.push((graph_id, c)), - _ => masters.push((graph_id, c)), - }, - Err(e) => { - counts.failed += 1; - ctx.logger - .warn(&format!("graph event {graph_id} convert failed: {e}")); + Ok(raw) => { + any_series |= is_series_master(&raw); + match convert_event(&raw, None) { + Ok(c) => match c.event_type { + EventType::Exception => exceptions.push((graph_id, c)), + _ => masters.push((graph_id, c)), + }, + Err(e) => { + counts.failed += 1; + ctx.logger + .warn(&format!("graph event {graph_id} convert failed: {e}")); + } } - }, + } Err(GraphError::Vanished) => counts.skipped += 1, Err(e) => { counts.failed += 1; @@ -101,6 +107,10 @@ pub fn reconcile_all( } } + if any_series { + exceptions.extend(fetch_calendar_exceptions(ctx, &cal.graph_id, counts)); + } + let mut master_by_graph_id: HashMap = masters.into_iter().collect(); for (ex_graph_id, ex) in exceptions { let Some(master_graph_id) = ex.series_master_id.clone() else { @@ -167,6 +177,73 @@ fn fetch_events( out } +fn is_series_master(raw: &Value) -> bool { + raw.get("recurrence").is_some_and(|r| !r.is_null()) +} + +fn exception_windows(years: i32) -> Vec<(String, String)> { + let today = chrono::Utc::now().date_naive(); + let span = chrono::Duration::days( + i64::from(years.max(1)) + .saturating_mul(365) + .min(api::EXCEPTION_WINDOW_MAX_DAYS), + ); + let fmt = |d: chrono::NaiveDate| d.format("%Y-%m-%dT00:00:00").to_string(); + vec![ + (fmt(today - span), fmt(today)), + (fmt(today), fmt(today + span)), + ] +} + +fn fetch_calendar_exceptions( + ctx: &GraphCoordinator<'_>, + calendar_id: &str, + counts: &mut TypeCounts, +) -> Vec<(String, ConvertedEvent)> { + let windows = exception_windows(ctx.exception_window_years); + let mut out = Vec::new(); + let mut seen: std::collections::HashSet = std::collections::HashSet::new(); + for (from, to) in windows { + let url = ctx + .endpoints + .calendar_exceptions(calendar_id, &from, &to, ctx.top); + match api::collect_all_values(ctx.client, &url, &[PREFER_TIMEZONE_UTC]) { + Ok(values) => { + if ctx.logger.enabled(LEVEL_PROGRESS) { + eprintln!( + "graph calendar {calendar_id} exceptions in {from}..{to}: {}", + values.len() + ); + } + for raw in values { + let Some(ex_id) = raw.get("id").and_then(Value::as_str).map(str::to_owned) + else { + continue; + }; + if !seen.insert(ex_id.clone()) { + continue; + } + match convert_event(&raw, None) { + Ok(c) => out.push((ex_id, c)), + Err(e) => { + counts.failed += 1; + ctx.logger + .warn(&format!("graph exception {ex_id} convert failed: {e}")); + } + } + } + } + Err(e) => { + counts.failed += 1; + ctx.logger.warn(&format!( + "graph calendar {calendar_id} exception lookup {from}..{to} failed: {e}" + )); + } + } + } + out +} + fn merge_exception_into(master: &mut ConvertedEvent, ex: &ConvertedEvent) { let Some(raw_key) = ex.original_start.as_deref() else { return; @@ -407,6 +484,48 @@ fn delete_vanished( mod tests { use super::*; + #[test] + fn exception_windows_never_exceed_the_graph_limit() { + for requested in [1, 5, 50] { + let windows = exception_windows(requested); + assert_eq!(windows.len(), 2, "one window behind today, one ahead"); + for (from, to) in windows { + let start = chrono::NaiveDate::parse_from_str(&from[..10], "%Y-%m-%d").unwrap(); + let end = chrono::NaiveDate::parse_from_str(&to[..10], "%Y-%m-%d").unwrap(); + let days = (end - start).num_days(); + assert!(days > 0, "window runs forwards: {from}..{to}"); + assert!( + days <= api::EXCEPTION_WINDOW_MAX_DAYS, + "requested {requested}y produced {days} days, over Graph's cap" + ); + } + } + } + + #[test] + fn cancelled_occurrence_keys_drop_the_oid_prefix() { + let raw = serde_json::json!({ + "id": "M", + "start": {"dateTime": "2026-09-21T11:00:00.0000000", "timeZone": "UTC"}, + "end": {"dateTime": "2026-09-21T11:30:00.0000000", "timeZone": "UTC"}, + "recurrence": { + "pattern": {"type": "daily", "interval": 1}, + "range": {"type": "numbered", "startDate": "2026-09-21", + "numberOfOccurrences": 6} + }, + "cancelledOccurrences": ["OID.AAMkAGI2TGuLAAA=.2026-09-24"] + }); + let converted = convert_event(&raw, None).unwrap(); + let overrides = converted.data["recurrenceOverrides"].as_object().unwrap(); + assert!( + overrides.contains_key("2026-09-24T11:00:00"), + "beta cancelledOccurrences are OID..; the key must be a LocalDateTime, \ + got {:?}", + overrides.keys().collect::>() + ); + assert_eq!(overrides["2026-09-24T11:00:00"]["excluded"], true); + } + #[test] fn override_key_utc_converts_to_master_local_time() { let local = normalise_override_key("2026-05-11T15:00:00Z", Some("America/New_York")); diff --git a/src/sync/import_exchange_graph/contacts.rs b/src/sync/import_exchange_graph/contacts.rs index fa0c5da..507b13d 100644 --- a/src/sync/import_exchange_graph/contacts.rs +++ b/src/sync/import_exchange_graph/contacts.rs @@ -59,10 +59,14 @@ pub fn reconcile_all( for id in &ids { server_total.insert(id.clone()); } - let new_ids: Vec = ids - .into_iter() - .filter(|id| !local.contains_key(id) && planned.insert(id.clone())) - .collect(); + let mut new_ids: Vec = Vec::new(); + for id in ids { + if local.contains_key(&id) { + counts.fetched += 1; + } else if planned.insert(id.clone()) { + new_ids.push(id); + } + } if new_ids.is_empty() { continue; } diff --git a/src/sync/import_exchange_graph/coordinator.rs b/src/sync/import_exchange_graph/coordinator.rs index 9e2fc5f..a101b84 100644 --- a/src/sync/import_exchange_graph/coordinator.rs +++ b/src/sync/import_exchange_graph/coordinator.rs @@ -42,6 +42,7 @@ pub struct GraphImportConfig { pub event_body_format: EventBodyFormat, pub graph_connections: usize, pub top: usize, + pub exception_window_years: i32, pub allow_source_change: bool, } @@ -55,6 +56,7 @@ pub struct GraphCoordinator<'a> { pub workers: usize, pub logger: crate::logging::Logger, pub event_body_format: EventBodyFormat, + pub exception_window_years: i32, } pub fn run(common: CommonConfig, config: GraphImportConfig) -> Result { @@ -129,6 +131,7 @@ pub fn run(common: CommonConfig, config: GraphImportConfig) -> Result Result Result = std::collections::HashSet::new(); let initial = endpoints.mail_folders_root(mailbox_kind, top); let level = crate::exchange_graph::api::collect_all_values(client, &initial, &[])?; - for f in &level { - if let Some(id) = f.get("id").and_then(Value::as_str) - && seen.insert(id.to_owned()) - { + for folder in level { + let Some(id) = folder.get("id").and_then(Value::as_str) else { + continue; + }; + if seen.insert(id.to_owned()) { frontier.push(id.to_owned()); + all.push(folder); } } - all.extend(level); while let Some(parent) = frontier.pop() { let url = endpoints.mail_folder_child_folders(&parent, top); let children = crate::exchange_graph::api::collect_all_values(client, &url, &[])?; - for c in &children { - if let Some(id) = c.get("id").and_then(Value::as_str) - && seen.insert(id.to_owned()) - { + for child in children { + let Some(id) = child.get("id").and_then(Value::as_str) else { + continue; + }; + if seen.insert(id.to_owned()) { frontier.push(id.to_owned()); + all.push(child); } } - all.extend(children); } Ok(all) } @@ -490,6 +502,7 @@ mod tests { event_body_format: EventBodyFormat::Text, graph_connections: 4, top: 100, + exception_window_years: 5, allow_source_change: false, }; let url = canonical_session_url(&config); diff --git a/src/sync/import_exchange_graph/files.rs b/src/sync/import_exchange_graph/files.rs new file mode 100644 index 0000000..dcc0a0b --- /dev/null +++ b/src/sync/import_exchange_graph/files.rs @@ -0,0 +1,438 @@ +/* + * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC + * + * SPDX-License-Identifier: Apache-2.0 OR MIT + */ + +use std::collections::HashMap; + +use rusqlite::{Connection, Transaction, params}; +use serde_json::Value; + +use crate::db::{blobs, exchange_graph_ids}; +use crate::error::Error; +use crate::exchange_graph::api; +use crate::exchange_graph::client::Accept; +use crate::exchange_graph::error::GraphError; +use crate::logging::LEVEL_PROGRESS; +use crate::sync::TypeCounts; +use crate::sync::import_jmap::pool::Pool; + +use super::coordinator::{CHUNK_SIZE, GraphCoordinator}; + +#[derive(Debug, Clone)] +struct DriveNode { + graph_id: String, + name: String, + is_folder: bool, + media_type: Option, + created: String, + modified: Option, + parent_graph_id: Option, +} + +pub fn reconcile_all( + conn: &mut Connection, + ctx: &GraphCoordinator<'_>, + counts: &mut TypeCounts, +) -> Result<(), Error> { + let Some(root_id) = drive_root_id(ctx) else { + ctx.logger + .warn("graph drive unavailable for this account; file import skipped"); + return Ok(()); + }; + + let nodes = match enumerate_drive(ctx, &root_id) { + Ok(v) => v, + Err(e) => { + counts.failed += 1; + ctx.logger + .warn(&format!("graph drive enumeration failed: {e}")); + return Ok(()); + } + }; + if ctx.logger.enabled(LEVEL_PROGRESS) { + eprintln!("graph drive items enumerated: {}", nodes.len()); + } + + let local: HashMap = + exchange_graph_ids::ids_of_type(conn, ctx.source_id, exchange_graph_ids::FILE_NODE)?; + let mut by_graph_id: HashMap = HashMap::new(); + let mut server_ids: Vec = Vec::new(); + + let (folders, files): (Vec, Vec) = + nodes.into_iter().partition(|n| n.is_folder); + + for node in order_by_depth(folders) { + server_ids.push(node.graph_id.clone()); + let parent_local = node + .parent_graph_id + .as_deref() + .and_then(|p| by_graph_id.get(p)) + .copied(); + let local_id = upsert_directory(conn, ctx, &node, parent_local, &local, counts)?; + by_graph_id.insert(node.graph_id.clone(), local_id); + } + + let mut wanted: Vec = Vec::new(); + for node in files { + server_ids.push(node.graph_id.clone()); + if local.contains_key(&node.graph_id) { + counts.fetched += 1; + } else { + wanted.push(node); + } + } + + download_and_insert(conn, ctx, &wanted, &by_graph_id, counts)?; + delete_vanished(conn, ctx.source_id, &local, &server_ids, counts, ctx)?; + Ok(()) +} + +fn drive_root_id(ctx: &GraphCoordinator<'_>) -> Option { + let value = ctx.client.get_json(&ctx.endpoints.drive_root()).ok()?; + value.get("id").and_then(Value::as_str).map(str::to_owned) +} + +fn enumerate_drive( + ctx: &GraphCoordinator<'_>, + root_id: &str, +) -> Result, GraphError> { + let mut out: Vec = Vec::new(); + let mut seen: std::collections::HashSet = std::collections::HashSet::new(); + let mut frontier: Vec = vec![root_id.to_owned()]; + seen.insert(root_id.to_owned()); + + while let Some(parent) = frontier.pop() { + let url = ctx.endpoints.drive_children(&parent, ctx.top); + let children = api::collect_all_values(ctx.client, &url, &[])?; + for child in &children { + let Some(node) = parse_node(child, &parent) else { + continue; + }; + if !seen.insert(node.graph_id.clone()) { + continue; + } + if node.is_folder { + frontier.push(node.graph_id.clone()); + } + out.push(node); + } + } + Ok(out) +} + +fn parse_node(value: &Value, parent_graph_id: &str) -> Option { + let graph_id = value.get("id").and_then(Value::as_str)?.to_owned(); + let name = value.get("name").and_then(Value::as_str)?.to_owned(); + let is_folder = value.get("folder").is_some_and(|f| !f.is_null()); + let is_file = value.get("file").is_some_and(|f| !f.is_null()); + if !is_folder && !is_file { + return None; + } + Some(DriveNode { + graph_id, + name, + is_folder, + media_type: value + .get("file") + .and_then(|f| f.get("mimeType")) + .and_then(Value::as_str) + .map(str::to_owned), + created: value + .get("createdDateTime") + .and_then(Value::as_str) + .unwrap_or("1970-01-01T00:00:00Z") + .to_owned(), + modified: value + .get("lastModifiedDateTime") + .and_then(Value::as_str) + .map(str::to_owned), + parent_graph_id: Some(parent_graph_id.to_owned()), + }) +} + +fn order_by_depth(folders: Vec) -> Vec { + let parents: HashMap> = folders + .iter() + .map(|n| (n.graph_id.clone(), n.parent_graph_id.clone())) + .collect(); + let depth = |start: &str| -> usize { + let mut d = 0; + let mut id = start.to_owned(); + let mut guard = 0; + while let Some(Some(parent)) = parents.get(&id) { + d += 1; + id = parent.clone(); + guard += 1; + if guard > 256 { + break; + } + } + d + }; + let mut ordered: Vec<(usize, DriveNode)> = folders + .iter() + .map(|n| (depth(n.graph_id.as_str()), n.clone())) + .collect(); + ordered.sort_by_key(|(d, _)| *d); + ordered.into_iter().map(|(_, n)| n).collect() +} + +fn upsert_directory( + conn: &mut Connection, + ctx: &GraphCoordinator<'_>, + node: &DriveNode, + parent_local: Option, + local: &HashMap, + counts: &mut TypeCounts, +) -> Result { + let tx = conn.unchecked_transaction()?; + let id = if let Some(existing) = local.get(&node.graph_id).copied() { + tx.execute( + "UPDATE file_nodes SET parent_id = ?1, name = ?2, modified = ?3 WHERE id = ?4", + params![parent_local, node.name, node.modified, existing], + )?; + counts.fetched += 1; + existing + } else { + tx.execute( + "INSERT INTO file_nodes (parent_id, node_type, blob_id, target, name, media_type, + created, modified, is_subscribed, role) + VALUES (?1, 'directory', NULL, NULL, ?2, NULL, ?3, ?4, 1, NULL)", + params![parent_local, node.name, node.created, node.modified], + )?; + let new_id = tx.last_insert_rowid(); + exchange_graph_ids::insert( + &tx, + ctx.source_id, + exchange_graph_ids::FILE_NODE, + &node.graph_id, + new_id, + )?; + counts.created += 1; + new_id + }; + tx.commit()?; + Ok(id) +} + +fn download_and_insert( + conn: &mut Connection, + ctx: &GraphCoordinator<'_>, + nodes: &[DriveNode], + by_graph_id: &HashMap, + counts: &mut TypeCounts, +) -> Result<(), Error> { + if nodes.is_empty() { + return Ok(()); + } + let node_by_id: HashMap<&str, &DriveNode> = + nodes.iter().map(|n| (n.graph_id.as_str(), n)).collect(); + let client = ctx.client.clone(); + let endpoints: crate::exchange_graph::api::Endpoints = (*ctx.endpoints).clone(); + type FetchResult = (String, Result, GraphError>); + let pool: Pool = Pool::new(ctx.workers, move |id: String| { + let url = endpoints.drive_item_content(&id); + match client.get_with_prefer(&url, Accept::Binary, &[]) { + Ok(resp) => (id, Ok(resp.body)), + Err(e) => (id, Err(e)), + } + }); + for node in nodes { + pool.submit(node.graph_id.clone()); + } + + let mut tx_opt: Option> = None; + let mut in_batch = 0usize; + for _ in 0..nodes.len() { + let Ok((graph_id, result)) = pool.results().recv() else { + break; + }; + let Some(node) = node_by_id.get(graph_id.as_str()).copied() else { + continue; + }; + match result { + Ok(bytes) => { + if tx_opt.is_none() { + tx_opt = Some(conn.unchecked_transaction()?); + } + let tx = tx_opt.as_mut().expect("tx is Some"); + let parent_local = node + .parent_graph_id + .as_deref() + .and_then(|p| by_graph_id.get(p)) + .copied(); + insert_file_in_tx(tx, ctx, node, parent_local, &bytes, counts)?; + in_batch += 1; + if in_batch >= CHUNK_SIZE { + if let Some(t) = tx_opt.take() { + t.commit()?; + } + in_batch = 0; + } + } + Err(GraphError::Vanished) => counts.skipped += 1, + Err(e) => { + counts.failed += 1; + ctx.logger + .warn(&format!("graph drive item {graph_id} download failed: {e}")); + } + } + } + if let Some(t) = tx_opt.take() { + t.commit()?; + } + Ok(()) +} + +fn insert_file_in_tx( + tx: &Transaction<'_>, + ctx: &GraphCoordinator<'_>, + node: &DriveNode, + parent_local: Option, + bytes: &[u8], + counts: &mut TypeCounts, +) -> Result<(), Error> { + let blob_id = blobs::intern_blob(tx, bytes)?; + tx.execute( + "INSERT INTO file_nodes (parent_id, node_type, blob_id, target, name, media_type, + created, modified, is_subscribed, role) + VALUES (?1, 'file', ?2, NULL, ?3, ?4, ?5, ?6, 1, NULL)", + params![ + parent_local, + blob_id, + node.name, + node.media_type, + node.created, + node.modified, + ], + )?; + let new_id = tx.last_insert_rowid(); + exchange_graph_ids::insert( + tx, + ctx.source_id, + exchange_graph_ids::FILE_NODE, + &node.graph_id, + new_id, + )?; + counts.created += 1; + Ok(()) +} + +fn delete_vanished( + conn: &mut Connection, + source_id: i64, + local: &HashMap, + server_ids: &[String], + counts: &mut TypeCounts, + ctx: &GraphCoordinator<'_>, +) -> Result<(), Error> { + let server: std::collections::HashSet<&str> = server_ids.iter().map(String::as_str).collect(); + let mut vanished: Vec<(&String, &i64)> = local + .iter() + .filter(|(graph_id, _)| !server.contains(graph_id.as_str())) + .collect(); + let depths = node_depths(conn)?; + vanished.sort_by_key(|(_, local_id)| std::cmp::Reverse(*depths.get(*local_id).unwrap_or(&0))); + + for (graph_id, local_id) in vanished { + let tx = conn.unchecked_transaction()?; + match tx.execute("DELETE FROM file_nodes WHERE id = ?1", params![local_id]) { + Ok(_) => { + exchange_graph_ids::delete( + &tx, + source_id, + exchange_graph_ids::FILE_NODE, + graph_id, + )?; + tx.commit()?; + counts.deleted += 1; + } + Err(e) => { + let _ = tx.rollback(); + ctx.logger.warn(&format!( + "file node {graph_id} (local id {local_id}) could not be deleted: {e}" + )); + counts.failed += 1; + } + } + } + Ok(()) +} + +fn node_depths(conn: &Connection) -> Result, Error> { + let mut parents: HashMap> = HashMap::new(); + { + let mut stmt = conn.prepare("SELECT id, parent_id FROM file_nodes")?; + let mut rows = stmt.query([])?; + while let Some(row) = rows.next()? { + parents.insert(row.get(0)?, row.get(1)?); + } + } + let mut depths = HashMap::new(); + for id in parents.keys().copied() { + let mut d = 0; + let mut cursor = id; + let mut seen: std::collections::HashSet = std::collections::HashSet::new(); + while seen.insert(cursor) { + match parents.get(&cursor).copied().flatten() { + Some(p) => { + d += 1; + cursor = p; + } + None => break, + } + } + depths.insert(id, d); + } + Ok(depths) +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + + #[test] + fn parse_node_rejects_items_without_file_or_folder_facet() { + let vault = json!({"id": "V", "name": "Personal Vault", "remoteItem": {}}); + assert!(parse_node(&vault, "root").is_none()); + } + + #[test] + fn parse_node_reads_folder_and_file_facets() { + let folder = json!({"id": "F", "name": "Docs", "folder": {"childCount": 2}}); + let parsed = parse_node(&folder, "root").expect("folder parses"); + assert!(parsed.is_folder); + + let file = json!({ + "id": "X", "name": "a.txt", "size": 12, + "file": {"mimeType": "text/plain"}, + "createdDateTime": "2026-01-01T00:00:00Z" + }); + let parsed = parse_node(&file, "root").expect("file parses"); + assert!(!parsed.is_folder); + assert_eq!(parsed.media_type.as_deref(), Some("text/plain")); + } + + #[test] + fn order_by_depth_puts_parents_before_children() { + let node = |id: &str, parent: Option<&str>| DriveNode { + graph_id: id.to_owned(), + name: id.to_owned(), + is_folder: true, + media_type: None, + created: "2026-01-01T00:00:00Z".to_owned(), + modified: None, + parent_graph_id: parent.map(str::to_owned), + }; + let ordered = order_by_depth(vec![ + node("C", Some("B")), + node("B", Some("A")), + node("A", None), + ]); + let names: Vec<&str> = ordered.iter().map(|n| n.graph_id.as_str()).collect(); + assert_eq!(names, ["A", "B", "C"]); + } +} diff --git a/src/sync/import_exchange_graph/folders.rs b/src/sync/import_exchange_graph/folders.rs index 37a4b2d..e9f6dce 100644 --- a/src/sync/import_exchange_graph/folders.rs +++ b/src/sync/import_exchange_graph/folders.rs @@ -16,6 +16,7 @@ use crate::exchange_graph::calendar_map::{graph_calendar_color_to_hex, windows_o use crate::exchange_graph::types::MailboxKind; use crate::logging::LEVEL_PROGRESS; use crate::sync::TypeCounts; +use crate::types::ObjectType; use super::coordinator::{CHUNK_SIZE, GraphCoordinator, enumerate_mail_folders}; @@ -242,18 +243,28 @@ pub fn reconcile_address_books( .filter_map(|f| f.get("id").and_then(Value::as_str)) .map(str::to_owned) .collect(); + let mut default_id: Option = None; + if let Some(default) = default_contact_folder(ctx) + && let Some(id) = default.get("id").and_then(Value::as_str) + { + default_id = Some(id.to_owned()); + if seen.insert(id.to_owned()) { + server.insert(0, default); + } + } let mut frontier: Vec = seen.iter().cloned().collect(); while let Some(parent) = frontier.pop() { let url = ctx.endpoints.contact_folder_children(&parent, ctx.top); let children = api::collect_all_values(ctx.client, &url, &[]).map_err(Error::from)?; - for c in &children { - if let Some(id) = c.get("id").and_then(Value::as_str) - && seen.insert(id.to_owned()) - { + for child in children { + let Some(id) = child.get("id").and_then(Value::as_str) else { + continue; + }; + if seen.insert(id.to_owned()) { frontier.push(id.to_owned()); + server.push(child); } } - server.extend(children); } let local: HashMap = @@ -272,19 +283,32 @@ pub fn reconcile_address_books( .and_then(Value::as_str) .unwrap_or("Contacts") .to_owned(); + let is_default = default_id.as_deref() == Some(graph_id); let existing = local.get(graph_id).copied(); let local_id = if let Some(id) = existing { + let is_default = crate::db::defaults::unique_default( + &tx, + ObjectType::AddressBook, + is_default, + Some(id), + )?; tx.execute( - "UPDATE address_books SET name = ?1 WHERE id = ?2", - params![name, id], + "UPDATE address_books SET name = ?1, is_default = ?2 WHERE id = ?3", + params![name, is_default as i64, id], )?; counts.fetched += 1; id } else { + let is_default = crate::db::defaults::unique_default( + &tx, + ObjectType::AddressBook, + is_default, + None, + )?; tx.execute( - "INSERT INTO address_books (name, sort_order, is_subscribed) - VALUES (?1, 0, 1)", - params![name], + "INSERT INTO address_books (name, sort_order, is_default, is_subscribed) + VALUES (?1, 0, ?2, 1)", + params![name, is_default as i64], )?; let new_id = tx.last_insert_rowid(); exchange_graph_ids::insert( @@ -317,6 +341,30 @@ pub fn reconcile_address_books( Ok(out) } +fn default_contact_folder(ctx: &GraphCoordinator<'_>) -> Option { + match ctx.client.get_json(&ctx.endpoints.default_contact_folder()) { + Ok(value) if value.get("id").and_then(Value::as_str).is_some() => return Some(value), + Ok(_) => {} + Err(e) => ctx.logger.warn(&format!( + "graph default contact folder lookup failed ({e}); \ + falling back to deriving it from an existing contact" + )), + } + let probe = ctx + .client + .get_json(&ctx.endpoints.any_contact_parent()) + .ok()?; + let parent = probe + .get("value") + .and_then(Value::as_array)? + .first()? + .get("parentFolderId") + .and_then(Value::as_str)?; + ctx.client + .get_json(&ctx.endpoints.contact_folder(parent)) + .ok() +} + fn mailbox_timezone(ctx: &GraphCoordinator<'_>) -> Option { let url = ctx.endpoints.mailbox_settings_timezone(); let value = ctx.client.get_json(&url).ok()?; diff --git a/src/sync/import_exchange_graph/messages.rs b/src/sync/import_exchange_graph/messages.rs index f23bbd9..a26fa34 100644 --- a/src/sync/import_exchange_graph/messages.rs +++ b/src/sync/import_exchange_graph/messages.rs @@ -5,7 +5,7 @@ */ use rusqlite::{Connection, Transaction, params}; -use serde_json::json; +use serde_json::{Value, json}; use crate::db::{blobs, exchange_graph_ids}; use crate::error::Error; @@ -35,7 +35,7 @@ pub fn reconcile_all( for folder in folders { let url = ctx.endpoints.folder_messages_ids(&folder.graph_id, ctx.top); - let ids = match api::collect_all_ids(ctx.client, &url, &[]) { + let stubs = match api::collect_all_values(ctx.client, &url, &[]) { Ok(v) => v, Err(e) => { ctx.logger.warn(&format!( @@ -51,16 +51,21 @@ pub fn reconcile_all( eprintln!( "graph folder {} enumerated {} messages", folder.graph_id, - ids.len() + stubs.len() ); } - for id in &ids { - server_total.insert(id.clone()); + let mut new_ids: Vec<(String, String)> = Vec::new(); + for stub in &stubs { + let Some(id) = stub.get("id").and_then(Value::as_str) else { + continue; + }; + server_total.insert(id.to_owned()); + if local.contains_key(id) { + counts.fetched += 1; + } else if planned.insert(id.to_owned()) { + new_ids.push((id.to_owned(), keyword_array(stub))); + } } - let new_ids: Vec = ids - .into_iter() - .filter(|id| !local.contains_key(id) && planned.insert(id.clone())) - .collect(); if new_ids.is_empty() { continue; } @@ -82,9 +87,13 @@ fn fetch_and_insert( conn: &mut Connection, ctx: &GraphCoordinator<'_>, folder: &MailFolder, - ids: &[String], + ids: &[(String, String)], counts: &mut TypeCounts, ) -> Result<(), Error> { + let keywords_by_id: std::collections::HashMap<&str, &str> = ids + .iter() + .map(|(id, kw)| (id.as_str(), kw.as_str())) + .collect(); let client = ctx.client.clone(); let endpoints: crate::exchange_graph::api::Endpoints = (*ctx.endpoints).clone(); type FetchResult = (String, Result, GraphError>); @@ -95,7 +104,7 @@ fn fetch_and_insert( Err(e) => (id, Err(e)), } }); - for id in ids { + for (id, _) in ids { pool.submit(id.clone()); } let mut tx_opt: Option> = None; @@ -110,7 +119,11 @@ fn fetch_and_insert( tx_opt = Some(conn.unchecked_transaction()?); } let tx = tx_opt.as_mut().expect("tx is Some"); - apply_message_in_tx(tx, ctx, folder, &graph_id, &bytes, counts)?; + let keywords = keywords_by_id + .get(graph_id.as_str()) + .copied() + .unwrap_or("[]"); + apply_message_in_tx(tx, ctx, folder, &graph_id, &bytes, keywords, counts)?; in_batch += 1; if in_batch >= CHUNK_SIZE { if let Some(t) = tx_opt.take() { @@ -141,13 +154,13 @@ fn apply_message_in_tx( folder: &MailFolder, graph_id: &str, bytes: &[u8], + keywords: &str, counts: &mut TypeCounts, ) -> Result<(), Error> { let (idx, date_header) = email_meta_from_blob(bytes); let message_match = index_to_json(&idx); let received_at = date_header.unwrap_or_else(|| "1970-01-01T00:00:00Z".to_owned()); let mailbox_ids = json!([folder.local_id]).to_string(); - let keywords = "[]".to_owned(); let blob_id = blobs::intern_blob(tx, bytes)?; tx.execute( @@ -167,6 +180,33 @@ fn apply_message_in_tx( Ok(()) } +fn keyword_array(stub: &Value) -> String { + let mut kws: Vec = Vec::new(); + if stub.get("isRead").and_then(Value::as_bool) == Some(true) { + kws.push("$seen".to_owned()); + } + if stub.get("isDraft").and_then(Value::as_bool) == Some(true) { + kws.push("$draft".to_owned()); + } + if stub.get("isReadReceiptRequested").and_then(Value::as_bool) == Some(true) { + kws.push("$notified".to_owned()); + } + if stub + .get("flag") + .and_then(|f| f.get("flagStatus")) + .and_then(Value::as_str) + == Some("flagged") + { + kws.push("$flagged".to_owned()); + } + if let Some(cats) = stub.get("categories").and_then(Value::as_array) { + for cat in cats.iter().filter_map(Value::as_str) { + kws.push(cat.to_ascii_lowercase()); + } + } + Value::Array(kws.into_iter().map(Value::String).collect()).to_string() +} + fn delete_vanished( conn: &mut Connection, source_id: i64, diff --git a/tests/mock_exchange_graph.rs b/tests/mock_exchange_graph.rs index c8ee9a2..1763e48 100644 --- a/tests/mock_exchange_graph.rs +++ b/tests/mock_exchange_graph.rs @@ -186,7 +186,7 @@ fn pagination_collects_ids_across_three_pages_with_one_empty() { .to_string(); let _m1 = server - .mock("GET", "/me/mailFolders/F1/messages?$top=100&$select=id") + .mock("GET", "/me/mailFolders/F1/messages?$top=100&$select=id,isRead,isDraft,isReadReceiptRequested,flag,categories") .with_status(200) .with_header("content-type", "application/json") .with_body(page1) @@ -220,7 +220,7 @@ fn pagination_terminates_when_no_next_link_present() { let base = server.url(); let body = json!({"value": [{"id": "A"}, {"id": "B"}]}).to_string(); let _m = server - .mock("GET", "/me/mailFolders/F/messages?$top=100&$select=id") + .mock("GET", "/me/mailFolders/F/messages?$top=100&$select=id,isRead,isDraft,isReadReceiptRequested,flag,categories") .with_status(200) .with_header("content-type", "application/json") .with_body(body) @@ -248,7 +248,7 @@ fn header_redelivery_is_asserted_via_mockito_match() { .to_string(); let page2 = json!({"value": [{"id": "OnlyOnPageTwo"}]}).to_string(); let _m1 = server - .mock("GET", "/me/mailFolders/F/messages?$top=100&$select=id") + .mock("GET", "/me/mailFolders/F/messages?$top=100&$select=id,isRead,isDraft,isReadReceiptRequested,flag,categories") .match_header("prefer", "IdType=\"ImmutableId\"") .match_header("authorization", "Bearer BEARER") .with_status(200) @@ -275,7 +275,7 @@ fn case_sensitive_immutable_id_round_trip() { let id = "AAkAAGFsaWNlAAA="; let body = json!({"value": [{"id": id}]}).to_string(); let _m = server - .mock("GET", "/me/mailFolders/F/messages?$top=100&$select=id") + .mock("GET", "/me/mailFolders/F/messages?$top=100&$select=id,isRead,isDraft,isReadReceiptRequested,flag,categories") .with_status(200) .with_header("content-type", "application/json") .with_body(body) @@ -589,6 +589,7 @@ fn integration_dry_run_against_mock_server_lists_three_surfaces() { event_body_format: vandelay::exchange_graph::types::EventBodyFormat::Text, graph_connections: 2, top: 100, + exception_window_years: 5, allow_source_change: false, }; let summary = vandelay::sync::import_exchange_graph::run(common, config).unwrap(); @@ -663,7 +664,7 @@ fn integration_full_run_mail_only_imports_mime_via_value() { }) .collect(); let _ids = server - .mock("GET", "/me/mailFolders/FMAIL/messages?$top=100&$select=id") + .mock("GET", "/me/mailFolders/FMAIL/messages?$top=100&$select=id,isRead,isDraft,isReadReceiptRequested,flag,categories") .with_status(200) .with_header("content-type", "application/json") .with_body(r#"{"value":[{"id":"MSG-1"}]}"#) @@ -698,6 +699,7 @@ fn integration_full_run_mail_only_imports_mime_via_value() { event_body_format: vandelay::exchange_graph::types::EventBodyFormat::Text, graph_connections: 2, top: 100, + exception_window_years: 5, allow_source_change: false, }; drop(tmp); @@ -759,7 +761,7 @@ fn integration_duplicate_message_id_does_not_abort_run() { }) .collect(); let _ids = server - .mock("GET", "/me/mailFolders/FMAIL/messages?$top=100&$select=id") + .mock("GET", "/me/mailFolders/FMAIL/messages?$top=100&$select=id,isRead,isDraft,isReadReceiptRequested,flag,categories") .with_status(200) .with_header("content-type", "application/json") .with_body(r#"{"value":[{"id":"MSG-1"},{"id":"MSG-1"}]}"#) @@ -794,6 +796,7 @@ fn integration_duplicate_message_id_does_not_abort_run() { event_body_format: vandelay::exchange_graph::types::EventBodyFormat::Text, graph_connections: 2, top: 100, + exception_window_years: 5, allow_source_change: false, }; drop(tmp); @@ -855,7 +858,7 @@ fn integration_full_run_is_convergent_on_second_invocation() { .create(); } let _ids = server - .mock("GET", "/me/mailFolders/FMAIL/messages?$top=100&$select=id") + .mock("GET", "/me/mailFolders/FMAIL/messages?$top=100&$select=id,isRead,isDraft,isReadReceiptRequested,flag,categories") .with_status(200) .with_header("content-type", "application/json") .with_body(r#"{"value":[{"id":"MSG-A"}]}"#) @@ -883,6 +886,7 @@ fn integration_full_run_is_convergent_on_second_invocation() { event_body_format: vandelay::exchange_graph::types::EventBodyFormat::Text, graph_connections: 2, top: 100, + exception_window_years: 5, allow_source_change: false, }; let make_common = |path: std::path::PathBuf| vandelay::sync::CommonConfig { @@ -962,6 +966,7 @@ fn source_change_protection_refuses_a_different_account() { event_body_format: vandelay::exchange_graph::types::EventBodyFormat::Text, graph_connections: 2, top: 100, + exception_window_years: 5, allow_source_change: false, }; let err = vandelay::sync::import_exchange_graph::run(common, config).unwrap_err(); @@ -1028,6 +1033,7 @@ fn make_config( event_body_format: vandelay::exchange_graph::types::EventBodyFormat::Text, graph_connections: 2, top: 100, + exception_window_years: 5, allow_source_change: false, } } @@ -1144,11 +1150,31 @@ fn series_master_with_exception_merges_into_recurrence_overrides() { .with_header("content-type", "application/json") .with_body( r#"{"value":[ - {"id":"MASTER","type":"seriesMaster","iCalUId":"uid-master"}, - {"id":"EX","type":"exception","seriesMasterId":"MASTER","iCalUId":"uid-master"} + {"id":"MASTER","type":"seriesMaster","iCalUId":"uid-master"} ]}"#, ) .create(); + server + .mock( + "GET", + Matcher::Regex(r"^/me/calendars/CAL1/calendarView\?startDateTime=".to_owned()), + ) + .with_status(200) + .with_header("content-type", "application/json") + .with_body( + r#"{"value":[{ + "id":"EX", + "iCalUId":"uid-master", + "type":"exception", + "seriesMasterId":"MASTER", + "originalStart":"2026-05-11T15:00:00Z", + "subject":"Moved sync", + "start":{"dateTime":"2026-05-11T16:00:00.0000000","timeZone":"UTC"}, + "end":{"dateTime":"2026-05-11T17:00:00.0000000","timeZone":"UTC"} + }]}"#, + ) + .expect_at_least(1) + .create(); server .mock("GET", "/me/events/MASTER") .with_status(200) @@ -1168,23 +1194,6 @@ fn series_master_with_exception_merges_into_recurrence_overrides() { }"#, ) .create(); - server - .mock("GET", "/me/events/EX") - .with_status(200) - .with_header("content-type", "application/json") - .with_body( - r#"{ - "id":"EX", - "iCalUId":"uid-master", - "type":"exception", - "seriesMasterId":"MASTER", - "originalStart":"2026-05-11T15:00:00Z", - "subject":"Moved sync", - "start":{"dateTime":"2026-05-11T16:00:00.0000000","timeZone":"UTC"}, - "end":{"dateTime":"2026-05-11T17:00:00.0000000","timeZone":"UTC"} - }"#, - ) - .create(); let base = server.url(); let archive = tempfile::NamedTempFile::new().unwrap().path().to_owned(); let summary = vandelay::sync::import_exchange_graph::run( @@ -1354,7 +1363,7 @@ fn hidden_mail_folder_has_is_subscribed_zero() { server .mock( "GET", - "/me/mailFolders/FVISIBLE/messages?$top=100&$select=id", + "/me/mailFolders/FVISIBLE/messages?$top=100&$select=id,isRead,isDraft,isReadReceiptRequested,flag,categories", ) .with_status(200) .with_header("content-type", "application/json") @@ -1363,7 +1372,7 @@ fn hidden_mail_folder_has_is_subscribed_zero() { server .mock( "GET", - "/me/mailFolders/FHIDDEN/messages?$top=100&$select=id", + "/me/mailFolders/FHIDDEN/messages?$top=100&$select=id,isRead,isDraft,isReadReceiptRequested,flag,categories", ) .with_status(200) .with_header("content-type", "application/json") @@ -1441,7 +1450,7 @@ fn well_known_folder_probes_assign_jmap_roles() { server .mock( "GET", - format!("/me/mailFolders/{fid}/messages?$top=100&$select=id").as_str(), + format!("/me/mailFolders/{fid}/messages?$top=100&$select=id,isRead,isDraft,isReadReceiptRequested,flag,categories").as_str(), ) .with_status(200) .with_header("content-type", "application/json") @@ -1515,7 +1524,7 @@ fn attachment_message_does_not_trigger_attachments_endpoint() { .create(); stub_well_known_folders(&mut server, "FMAIL"); server - .mock("GET", "/me/mailFolders/FMAIL/messages?$top=100&$select=id") + .mock("GET", "/me/mailFolders/FMAIL/messages?$top=100&$select=id,isRead,isDraft,isReadReceiptRequested,flag,categories") .with_status(200) .with_header("content-type", "application/json") .with_body(r#"{"value":[{"id":"MSG-ATT"}]}"#) @@ -1708,7 +1717,7 @@ fn full_run_records_graph_id_in_sync_id_exchange_graph_with_padding() { let _ids = server .mock( "GET", - "/me/mailFolders/AAkA-Padded%3D%3D/messages?$top=100&$select=id", + "/me/mailFolders/AAkA-Padded%3D%3D/messages?$top=100&$select=id,isRead,isDraft,isReadReceiptRequested,flag,categories", ) .with_status(200) .with_header("content-type", "application/json") @@ -1736,6 +1745,7 @@ fn full_run_records_graph_id_in_sync_id_exchange_graph_with_padding() { event_body_format: vandelay::exchange_graph::types::EventBodyFormat::Text, graph_connections: 2, top: 100, + exception_window_years: 5, allow_source_change: false, }; drop(tmp); @@ -1806,6 +1816,7 @@ fn archive_mailbox_kind_encodes_synthetic_suffix_in_account_id() { event_body_format: vandelay::exchange_graph::types::EventBodyFormat::Text, graph_connections: 2, top: 100, + exception_window_years: 5, allow_source_change: false, }; drop(tmp); @@ -1871,6 +1882,7 @@ fn allow_source_change_permits_overwriting_a_different_account() { event_body_format: vandelay::exchange_graph::types::EventBodyFormat::Text, graph_connections: 2, top: 100, + exception_window_years: 5, allow_source_change: true, }; let result = vandelay::sync::import_exchange_graph::run(common, config); @@ -1952,6 +1964,7 @@ fn dry_run_makes_no_per_item_get_and_no_sqlite_writes() { event_body_format: vandelay::exchange_graph::types::EventBodyFormat::Text, graph_connections: 2, top: 100, + exception_window_years: 5, allow_source_change: false, }; drop(tmp); @@ -2096,7 +2109,7 @@ fn folder_enumeration_failure_skips_vanished_deletion() { .create(); } let _enum_fail = server - .mock("GET", "/me/mailFolders/FMAIL/messages?$top=100&$select=id") + .mock("GET", "/me/mailFolders/FMAIL/messages?$top=100&$select=id,isRead,isDraft,isReadReceiptRequested,flag,categories") .with_status(503) .with_body("transient outage") .create(); @@ -2120,6 +2133,7 @@ fn folder_enumeration_failure_skips_vanished_deletion() { event_body_format: vandelay::exchange_graph::types::EventBodyFormat::Text, graph_connections: 2, top: 100, + exception_window_years: 5, allow_source_change: true, }; let _ = vandelay::sync::import_exchange_graph::run(common, config).unwrap(); @@ -2147,7 +2161,14 @@ fn stub_contacts_surface(server: &mut Server) { .mock("GET", "/me/contactFolders?$top=100") .with_status(200) .with_header("content-type", "application/json") - .with_body(r#"{"value":[{"id":"CON1","displayName":"Contacts"}]}"#) + .with_body(r#"{"value":[]}"#) + .expect_at_least(0) + .create(); + server + .mock("GET", "/me/contactFolders/contacts") + .with_status(200) + .with_header("content-type", "application/json") + .with_body(r#"{"id":"CON1","displayName":"Contacts","wellKnownName":"contacts"}"#) .expect_at_least(0) .create(); server @@ -2199,7 +2220,7 @@ fn stub_mail_surface(server: &mut Server) { .create(); stub_well_known_folders(server, "FMAIL"); server - .mock("GET", "/me/mailFolders/FMAIL/messages?$top=100&$select=id") + .mock("GET", "/me/mailFolders/FMAIL/messages?$top=100&$select=id,isRead,isDraft,isReadReceiptRequested,flag,categories") .with_status(200) .with_header("content-type", "application/json") .with_body(r#"{"value":[{"id":"MSG-S"}]}"#) @@ -2342,3 +2363,375 @@ fn mail_surface_imports_only_mailboxes_and_emails() { assert_eq!(row_count(&conn, "calendars"), 0); assert_eq!(row_count(&conn, "calendar_events"), 0); } + +fn stub_principal(server: &mut Server) { + server + .mock("GET", Matcher::Regex(r"^/me\?\$select=id".to_owned())) + .with_status(200) + .with_header("content-type", "application/json") + .with_body(r#"{"id":"uid","userPrincipalName":"alice@x.com"}"#) + .expect_at_least(0) + .create(); +} + +fn json_mock(server: &mut Server, path: &str, body: &str) { + server + .mock("GET", path) + .with_status(200) + .with_header("content-type", "application/json") + .with_body(body) + .expect_at_least(0) + .create(); +} + +#[test] +fn default_contact_folder_is_imported_although_contactfolders_omits_it() { + let mut server = Server::new(); + stub_principal(&mut server); + json_mock( + &mut server, + "/me/contactFolders?$top=100", + r#"{"value":[]}"#, + ); + json_mock( + &mut server, + "/me/contactFolders/contacts", + r#"{"id":"DEFAULT","displayName":"Contacts","wellKnownName":"contacts"}"#, + ); + json_mock( + &mut server, + "/me/contactFolders/DEFAULT/childFolders?$top=100", + r#"{"value":[]}"#, + ); + json_mock( + &mut server, + "/me/contactFolders/DEFAULT/contacts?$top=100&$select=id", + r#"{"value":[{"id":"C1"},{"id":"C2"}]}"#, + ); + json_mock( + &mut server, + "/me/contacts/C1", + r#"{"id":"C1","displayName":"Alice"}"#, + ); + json_mock( + &mut server, + "/me/contacts/C2", + r#"{"id":"C2","displayName":"Bob"}"#, + ); + + let base = server.url(); + let archive = tempfile::NamedTempFile::new().unwrap().path().to_owned(); + let summary = vandelay::sync::import_exchange_graph::run( + make_common(archive.clone()), + make_config(base, None, surfaces("contacts")), + ) + .unwrap(); + + let count = |name: &str| { + summary + .per_type + .iter() + .find(|(t, _)| *t == name) + .map(|(_, c)| c.created) + .unwrap_or(0) + }; + assert_eq!( + count("addressbook"), + 1, + "the default Contacts folder must become an address book even though \ + /me/contactFolders never lists it" + ); + assert_eq!( + count("contactcard"), + 2, + "both default-folder contacts import" + ); + + let conn = vandelay::db::init::open(&archive).unwrap(); + let is_default: i64 = conn + .query_row("SELECT is_default FROM address_books", [], |row| row.get(0)) + .unwrap(); + assert_eq!(is_default, 1, "the well-known folder is the default book"); +} + +#[test] +fn default_contact_folder_falls_back_to_parent_of_an_existing_contact() { + let mut server = Server::new(); + stub_principal(&mut server); + json_mock( + &mut server, + "/me/contactFolders?$top=100", + r#"{"value":[]}"#, + ); + server + .mock("GET", "/me/contactFolders/contacts") + .with_status(404) + .with_body(r#"{"error":{"code":"ErrorItemNotFound"}}"#) + .expect_at_least(1) + .create(); + json_mock( + &mut server, + "/me/contacts?$top=1&$select=id,parentFolderId", + r#"{"value":[{"id":"C1","parentFolderId":"DERIVED"}]}"#, + ); + json_mock( + &mut server, + "/me/contactFolders/DERIVED", + r#"{"id":"DERIVED","displayName":"Kontakte"}"#, + ); + json_mock( + &mut server, + "/me/contactFolders/DERIVED/childFolders?$top=100", + r#"{"value":[]}"#, + ); + json_mock( + &mut server, + "/me/contactFolders/DERIVED/contacts?$top=100&$select=id", + r#"{"value":[{"id":"C1"}]}"#, + ); + json_mock( + &mut server, + "/me/contacts/C1", + r#"{"id":"C1","displayName":"Alice"}"#, + ); + + let base = server.url(); + let archive = tempfile::NamedTempFile::new().unwrap().path().to_owned(); + vandelay::sync::import_exchange_graph::run( + make_common(archive.clone()), + make_config(base, None, surfaces("contacts")), + ) + .unwrap(); + + let conn = vandelay::db::init::open(&archive).unwrap(); + let name: String = conn + .query_row("SELECT name FROM address_books", [], |row| row.get(0)) + .unwrap(); + assert_eq!( + name, "Kontakte", + "when the well-known name is unavailable the folder is derived from a contact's parent" + ); +} + +#[test] +fn contact_folder_reachable_by_two_paths_is_inserted_once() { + let mut server = Server::new(); + stub_principal(&mut server); + json_mock( + &mut server, + "/me/contactFolders?$top=100", + r#"{"value":[{"id":"CHILD","displayName":"Work","parentFolderId":"DEFAULT"}]}"#, + ); + json_mock( + &mut server, + "/me/contactFolders/contacts", + r#"{"id":"DEFAULT","displayName":"Contacts","wellKnownName":"contacts"}"#, + ); + json_mock( + &mut server, + "/me/contactFolders/DEFAULT/childFolders?$top=100", + r#"{"value":[{"id":"CHILD","displayName":"Work","parentFolderId":"DEFAULT"}]}"#, + ); + json_mock( + &mut server, + "/me/contactFolders/CHILD/childFolders?$top=100", + r#"{"value":[]}"#, + ); + for folder in ["DEFAULT", "CHILD"] { + json_mock( + &mut server, + &format!("/me/contactFolders/{folder}/contacts?$top=100&$select=id"), + r#"{"value":[]}"#, + ); + } + + let base = server.url(); + let archive = tempfile::NamedTempFile::new().unwrap().path().to_owned(); + let summary = vandelay::sync::import_exchange_graph::run( + make_common(archive.clone()), + make_config(base, None, surfaces("contacts")), + ) + .expect("a folder reachable by two paths must not violate the id-mapping uniqueness"); + + let created = summary + .per_type + .iter() + .find(|(t, _)| *t == "addressbook") + .map(|(_, c)| c.created) + .unwrap_or(0); + assert_eq!( + created, 2, + "the default folder and its one child, each once" + ); +} + +#[test] +fn message_state_becomes_jmap_keywords() { + let mut server = Server::new(); + stub_principal(&mut server); + json_mock( + &mut server, + "/me/mailFolders?$top=100&includeHiddenFolders=true", + r#"{"value":[{"id":"F1","displayName":"Inbox","isHidden":false}]}"#, + ); + json_mock( + &mut server, + "/me/mailFolders/F1/childFolders?$top=100&includeHiddenFolders=true", + r#"{"value":[]}"#, + ); + for name in [ + "inbox", + "drafts", + "sentitems", + "deleteditems", + "junkemail", + "archive", + ] { + server + .mock("GET", format!("/me/mailFolders/{name}?$select=id").as_str()) + .with_status(404) + .expect_at_least(0) + .create(); + } + json_mock( + &mut server, + "/me/mailFolders/F1/messages?$top=100&$select=id,isRead,isDraft,isReadReceiptRequested,flag,categories", + r#"{"value":[ + {"id":"M1","isRead":true,"isDraft":false,"isReadReceiptRequested":true, + "flag":{"flagStatus":"flagged"},"categories":["Red Category","VIP"]}, + {"id":"M2","isRead":false,"isDraft":true, + "flag":{"flagStatus":"notFlagged"},"categories":[]} + ]}"#, + ); + for id in ["M1", "M2"] { + server + .mock("GET", format!("/me/messages/{id}/$value").as_str()) + .with_status(200) + .with_header("content-type", "text/plain") + .with_body(format!( + "From: a@x.com\r\nTo: b@x.com\r\nSubject: {id}\r\nMessage-ID: <{id}@x>\r\n\r\nbody\r\n" + )) + .expect_at_least(0) + .create(); + } + + let base = server.url(); + let archive = tempfile::NamedTempFile::new().unwrap().path().to_owned(); + vandelay::sync::import_exchange_graph::run( + make_common(archive.clone()), + make_config(base, None, surfaces("mail")), + ) + .unwrap(); + + let conn = vandelay::db::init::open(&archive).unwrap(); + let mut stmt = conn + .prepare("SELECT keywords FROM emails ORDER BY id") + .unwrap(); + let rows: Vec = stmt + .query_map([], |row| row.get::<_, String>(0)) + .unwrap() + .map(|r| r.unwrap()) + .collect(); + let all = rows.join(" "); + for expected in [ + "$seen", + "$notified", + "$flagged", + "red category", + "vip", + "$draft", + ] { + assert!(all.contains(expected), "missing {expected} in {all}"); + } +} + +#[test] +fn drive_items_import_as_file_nodes_and_skip_facetless_items() { + let mut server = Server::new(); + stub_principal(&mut server); + json_mock( + &mut server, + "/me/drive/root?$select=id,name", + r#"{"id":"ROOT","name":"root"}"#, + ); + let select = "id,name,size,folder,file,package,remoteItem,createdDateTime,lastModifiedDateTime"; + json_mock( + &mut server, + &format!("/me/drive/items/ROOT/children?$top=100&$select={select}"), + r#"{"value":[ + {"id":"D1","name":"Docs","folder":{"childCount":1}, + "createdDateTime":"2026-01-01T00:00:00Z","lastModifiedDateTime":"2026-01-02T00:00:00Z"}, + {"id":"V1","name":"Personal Vault","remoteItem":{}, + "createdDateTime":"2026-01-01T00:00:00Z"}, + {"id":"F1","name":"top.txt","size":5,"file":{"mimeType":"text/plain"}, + "createdDateTime":"2026-01-01T00:00:00Z","lastModifiedDateTime":"2026-01-02T00:00:00Z"} + ]}"#, + ); + json_mock( + &mut server, + &format!("/me/drive/items/D1/children?$top=100&$select={select}"), + r#"{"value":[ + {"id":"F2","name":"inner.bin","size":3,"file":{"mimeType":"application/octet-stream"}, + "createdDateTime":"2026-01-01T00:00:00Z"} + ]}"#, + ); + for (id, body) in [("F1", "hello"), ("F2", "abc")] { + server + .mock("GET", format!("/me/drive/items/{id}/content").as_str()) + .with_status(200) + .with_header("content-type", "application/octet-stream") + .with_body(body) + .expect_at_least(0) + .create(); + } + + let base = server.url(); + let archive = tempfile::NamedTempFile::new().unwrap().path().to_owned(); + let summary = vandelay::sync::import_exchange_graph::run( + make_common(archive.clone()), + make_config(base, None, surfaces("files")), + ) + .unwrap(); + + let created = summary + .per_type + .iter() + .find(|(t, _)| *t == "filenode") + .map(|(_, c)| c.created) + .unwrap_or(0); + assert_eq!( + created, 3, + "one directory and two files; the remoteItem has no file or folder facet and is skipped" + ); + + let conn = vandelay::db::init::open(&archive).unwrap(); + let vault: i64 = conn + .query_row( + "SELECT count(*) FROM file_nodes WHERE name = 'Personal Vault'", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(vault, 0, "a facetless drive item must never become a node"); + + let nested: i64 = conn + .query_row( + "SELECT count(*) FROM file_nodes child JOIN file_nodes parent + ON child.parent_id = parent.id + WHERE child.name = 'inner.bin' AND parent.name = 'Docs'", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(nested, 1, "child files hang off their directory"); + + let body: Vec = conn + .query_row( + "SELECT b.data FROM file_nodes f JOIN blobs b ON b.id = f.blob_id + WHERE f.name = 'top.txt'", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(body, b"hello", "file content is stored verbatim"); +}