diff --git a/README.md b/README.md index 3a110c4..ffa555e 100644 --- a/README.md +++ b/README.md @@ -57,6 +57,56 @@ Because the archive is a self-contained SQLite file that fully describes one acc - **Source-change protection:** An archive remembers which account it was filled from; pointing it at a different one fails unless explicitly permitted. - **Read-only inspection:** A built-in `inspect` command dumps any object type from an archive for verification. +## Install + +```sh +# macOS / Linux +curl --proto '=https' --tlsv1.2 -LsSf \ + https://github.com/stalwartlabs/vandelay/releases/latest/download/vandelay-installer.sh | sh + +# Homebrew +brew install stalwartlabs/tap/vandelay + +# Windows +powershell -ExecutionPolicy Bypass -c "irm https://github.com/stalwartlabs/vandelay/releases/latest/download/vandelay-installer.ps1 | iex" + +# npm +npm install -g @stalwartlabs/vandelay + +# From source +cargo install --path . +``` + +A signed `.msi` is also published with each release. + +## Quick start + +A typical run is two commands, one to capture a source account into a local SQLite archive and one to push that archive into a JMAP target. + +```sh +# 1. Import an IMAP mailbox into a fresh archive. +export VANDELAY_PASSWORD='source-app-password' +vandelay import imap \ + --url imaps://imap.example.com \ + --auth-basic alice@example.com \ + alice.sqlite + +# 2. Peek at what landed. +vandelay inspect alice.sqlite # per-type summary +vandelay inspect alice.sqlite mailbox # mailbox tree +vandelay inspect alice.sqlite email --limit 5 + +# 3. Push the archive into a target JMAP server. +export VANDELAY_PASSWORD='target-password' +vandelay export \ + --url https://jmap.example.org \ + --auth-basic alice@example.org \ + --account-name alice@example.org \ + alice.sqlite +``` + +Both commands are convergent: rerun either to resume an interrupted run, or rerun `import` later to pick up new mail since the last snapshot. Use `--dry-run` on either side to compute the full plan without writing. + ## CLI quick reference ``` diff --git a/src/dav/discover.rs b/src/dav/discover.rs index f6620fb..eab0d4f 100644 --- a/src/dav/discover.rs +++ b/src/dav/discover.rs @@ -136,13 +136,43 @@ fn try_treat_url_as_home_or_collection( if collections.is_empty() { return Ok(None); } + let principal_url = resolve_principal_url(client, kind, url)?; Ok(Some(Discovery { - principal_url: None, + principal_url, home_set_url: url.to_owned(), collections, })) } +fn resolve_principal_url( + client: &DavClient, + kind: DavKind, + url: &str, +) -> Result, DiscoveryError> { + if matches!(kind, DavKind::Webdav) { + return Ok(None); + } + let body = xml::propfind_current_user_principal(); + let ms = match client.propfind_responses(url, 0, &body, url) { + Ok(ms) => ms, + Err(JmapError::HttpStatus { status, .. }) if (400..600).contains(&status) => { + return Ok(None); + } + Err(JmapError::RetriesExhausted(_)) => return Ok(None), + Err(e) => return Err(DiscoveryError::Transport(e)), + }; + let final_url = ms.final_url.clone(); + let Some(principal_href) = ms + .responses + .iter() + .find_map(|r| r.props.current_user_principal.as_deref()) + else { + return Ok(None); + }; + let resolved = absolute_url(&final_url, &normalise(&final_url, principal_href)?)?; + Ok(Some(resolved)) +} + fn try_via_principal( client: &DavClient, kind: DavKind, diff --git a/src/main.rs b/src/main.rs index 26eea8e..c6d092c 100644 --- a/src/main.rs +++ b/src/main.rs @@ -98,8 +98,8 @@ fn fail(err: &Error) -> i32 { fn report(summary: &Summary) { for (type_name, counts) in &summary.per_type { println!( - "{type_name}: created={} fetched={} deleted={} skipped={} failed={}", - counts.created, counts.fetched, counts.deleted, counts.skipped, counts.failed + "{type_name}: created={} fetched={} updated={} deleted={} skipped={} failed={}", + counts.created, counts.fetched, counts.updated, counts.deleted, counts.skipped, counts.failed ); } } diff --git a/src/sync/export/sieve.rs b/src/sync/export/sieve.rs index 1cdb59b..cbb88dc 100644 --- a/src/sync/export/sieve.rs +++ b/src/sync/export/sieve.rs @@ -4,19 +4,16 @@ * SPDX-License-Identifier: Apache-2.0 OR MIT */ -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use serde_json::{Value, json}; use super::common::{create_batch, jid, target_get_all}; use super::{Maps, Net, Plan, Uploader}; -use crate::db; use crate::error::Error; -use crate::jmap::blobxfer; use crate::jmap::request::Request; use crate::logging::Logger; use crate::sync::import_jmap::mapping::{SIEVE_SELECT, row_to_sieve_script}; -use crate::sync::keys::blake3_bytes; use crate::sync::{Context, TypeCounts}; use crate::types::ObjectType; @@ -30,21 +27,12 @@ pub fn reconcile( let ty = ObjectType::SieveScript; let targets = target_get_all(net, ty).map_err(Error::from)?; - let mut target_by_key: HashMap<[u8; 32], String> = HashMap::new(); + let mut target_by_name: HashMap = HashMap::new(); for t in &targets { - let (Some(id), Some(blob)) = (jid(t), t.get("blobId").and_then(Value::as_str)) else { + let (Some(id), Some(name)) = (jid(t), t.get("name").and_then(Value::as_str)) else { continue; }; - let bytes = blobxfer::download_bytes( - &net.client, - &net.session, - &net.account, - blob, - "application/sieve", - "script", - ) - .map_err(Error::from)?; - target_by_key.insert(blake3_bytes(&bytes), id); + target_by_name.insert(name.to_owned(), id); } let locals: Vec<(i64, Option, bool, i64)> = { @@ -71,13 +59,10 @@ pub fn reconcile( let mut uploader = Uploader::new(net, &ctx.conn); for (local, name, is_active, blob_local) in &locals { - let bytes = db::blobs::blob_bytes(&ctx.conn, *blob_local) - .map_err(|e| Error::Partial(e.to_string()))? - .ok_or_else(|| Error::Partial("sieve blob missing".to_owned()))?; - let key = blake3_bytes(&bytes); - let target_id = if let Some(id) = target_by_key.get(&key) { + let matched = name.as_ref().and_then(|n| target_by_name.get(n)).cloned(); + let target_id = if let Some(id) = matched { counts.skipped += 1; - id.clone() + id } else { let blob_id = uploader .upload_with(*blob_local, "application/sieve") @@ -92,7 +77,9 @@ pub fn reconcile( match outcome.created.first().and_then(|(_, v)| jid(v)) { Some(id) => { counts.created += 1; - target_by_key.insert(key, id.clone()); + if let Some(n) = name { + target_by_name.insert(n.clone(), id.clone()); + } id } None => { @@ -128,18 +115,10 @@ pub fn reconcile( } } - let local_keys: std::collections::HashSet<[u8; 32]> = locals + let local_names: HashSet = locals.iter().filter_map(|(_, n, _, _)| n.clone()).collect(); + let mut prune_candidates: Vec = target_by_name .iter() - .filter_map(|(_, _, _, b)| { - db::blobs::blob_bytes(&ctx.conn, *b) - .ok() - .flatten() - .map(|by| blake3_bytes(&by)) - }) - .collect(); - let mut prune_candidates: Vec = target_by_key - .iter() - .filter(|(k, _)| !local_keys.contains(*k)) + .filter(|(name, _)| !local_names.contains(*name)) .map(|(_, id)| id.clone()) .collect(); prune_candidates.sort(); diff --git a/src/sync/import_imap/coordinator.rs b/src/sync/import_imap/coordinator.rs index d0d96a8..427ee3c 100644 --- a/src/sync/import_imap/coordinator.rs +++ b/src/sync/import_imap/coordinator.rs @@ -826,6 +826,18 @@ fn reconcile_folder( tx.commit()?; } + if !diff.present.is_empty() { + counts.updated += refresh_present_flags( + conn, + client, + control_ctx, + &folder.name, + uidvalidity, + &diff.present, + opts, + )?; + } + let now = OffsetDateTime::now_utc() .format(&Rfc3339) .unwrap_or_else(|_| String::from("1970-01-01T00:00:00Z")); @@ -961,6 +973,93 @@ fn insert_single_message( Ok(()) } +fn refresh_present_flags( + conn: &mut Connection, + client: &mut ImapClient, + control_ctx: &ControlCtx, + folder: &str, + uidvalidity: u32, + present: &[u32], + opts: RunOpts, +) -> Result { + let RunOpts { + source_id, + fetch_batch, + include_deleted, + logger: _, + } = opts; + let mut updated: u64 = 0; + let tx = conn.transaction()?; + for batch in chunks(present, fetch_batch.max(1)) { + let set = command::format_uid_set(batch, true); + let resp = control_run_collect( + client, + control_ctx, + &command::uid_fetch(&set, &["UID", "FLAGS"]), + )?; + for u in &resp.untagged { + let Some(attrs) = fetch::extract(u) else { + continue; + }; + let Some(uid) = attrs.uid else { + continue; + }; + let translation = translate_flags(&attrs.flags, include_deleted); + if translation.has_deleted_flag && !include_deleted { + continue; + } + let Some((local_id, existing)) = + load_email_keywords(&tx, source_id, folder, uidvalidity, uid)? + else { + continue; + }; + if keyword_set_differs(&existing, &translation.keywords) { + let json = Value::Array( + translation + .keywords + .iter() + .cloned() + .map(Value::String) + .collect(), + ); + tx.execute( + "UPDATE emails SET keywords = ?1 WHERE id = ?2", + params![json.to_string(), local_id], + )?; + updated += 1; + } + } + } + tx.commit()?; + Ok(updated) +} + +fn load_email_keywords( + tx: &rusqlite::Transaction<'_>, + source_id: i64, + folder: &str, + uidvalidity: u32, + uid: u32, +) -> Result)>, Error> { + let Some(local_id) = db::imap_ids::local_for_email(tx, source_id, folder, uidvalidity, uid)? + else { + return Ok(None); + }; + let kw_json: String = tx.query_row( + "SELECT keywords FROM emails WHERE id = ?1", + params![local_id], + |r| r.get(0), + )?; + let kws: Vec = serde_json::from_str(&kw_json).unwrap_or_default(); + Ok(Some((local_id, kws))) +} + +fn keyword_set_differs(a: &[String], b: &[String]) -> bool { + let sa: HashSet<&str> = a.iter().map(String::as_str).collect(); + let sb: HashSet<&str> = b.iter().map(String::as_str).collect(); + sa != sb +} + fn dry_run_summary( conn: &Connection, source_id: Option, diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 1efbb14..31bbe73 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -75,6 +75,7 @@ pub struct ExportConfig { pub struct TypeCounts { pub created: u64, pub fetched: u64, + pub updated: u64, pub deleted: u64, pub skipped: u64, pub failed: u64, diff --git a/tests/mock_dav.rs b/tests/mock_dav.rs index 786f63e..f2cfe94 100644 --- a/tests/mock_dav.rs +++ b/tests/mock_dav.rs @@ -66,6 +66,66 @@ fn discovery_uses_url_as_homeset_when_collection_present() { assert!(disc.collections[0].props.is_calendar); } +#[test] +fn discovery_resolves_per_user_principal_even_when_url_lists_collections() { + let mut server = mockito::Server::new(); + let url = server.url(); + + let collections_body = format!( + r#" + + + {url}/dav/cal/secondary@vandelay.org/default/ + + + + Default + + HTTP/1.1 200 OK + + +"# + ); + let _listing = server + .mock("PROPFIND", "/dav/cal/") + .match_header("depth", "1") + .with_status(207) + .with_header("content-type", "application/xml; charset=utf-8") + .with_body(&collections_body) + .create(); + + let principal_body = format!( + r#" + + + {url}/dav/cal/ + + + {url}/dav/principals/secondary@vandelay.org/ + + HTTP/1.1 200 OK + + +"# + ); + let _principal = server + .mock("PROPFIND", "/dav/cal/") + .match_header("depth", "0") + .with_status(207) + .with_header("content-type", "application/xml; charset=utf-8") + .with_body(&principal_body) + .create(); + + let c = client(0); + let disc = discover(&c, DavKind::Caldav, &format!("{url}/dav/cal/")).expect("discover"); + assert_eq!(disc.collections.len(), 1); + assert_eq!( + disc.principal_url.as_deref(), + Some(format!("{url}/dav/principals/secondary@vandelay.org/").as_str()), + "account identity must be the per-user principal, not the shared base DAV root" + ); +} + #[test] fn discovery_falls_through_principal_to_home_set() { let mut server = mockito::Server::new(); @@ -1129,6 +1189,133 @@ fn dry_run_writes_nothing_but_emits_per_collection_counts() { assert_eq!(event_counts.1.created, 2); } +#[test] +fn dav_source_change_protection_fires_across_users_on_same_root() { + use base64::Engine; + use base64::engine::general_purpose::STANDARD; + use std::path::PathBuf; + use vandelay::logging::Logger; + use vandelay::sync::CommonConfig; + use vandelay::sync::import_dav::{DavAuth, DavImportConfig, DavKindArg, run}; + + let mut server = mockito::Server::new(); + let url = server.url(); + + let listing = format!( + r#" + + + {url}/dav/cal/shared/default/ + + + + Default + + HTTP/1.1 200 OK + + +"# + ); + let _listing = server + .mock("PROPFIND", "/dav/cal/") + .match_header("depth", "1") + .with_status(207) + .with_header("content-type", "application/xml; charset=utf-8") + .with_body(&listing) + .create(); + + let empty_items = format!( + r#" + + + {url}/dav/cal/shared/default/ + + + HTTP/1.1 200 OK + + +"# + ); + let _items = server + .mock("PROPFIND", "/dav/cal/shared/default/") + .with_status(207) + .with_header("content-type", "application/xml; charset=utf-8") + .with_body(&empty_items) + .create(); + + let auth_a = format!("Basic {}", STANDARD.encode("secondary@vandelay.org:passA")); + let auth_b = format!("Basic {}", STANDARD.encode("tertiary@vandelay.org:passB")); + let principal = |who: &str| { + format!( + r#" + + + {url}/dav/cal/ + + {url}/dav/principals/{who}/ + HTTP/1.1 200 OK + + +"# + ) + }; + let _pa = server + .mock("PROPFIND", "/dav/cal/") + .match_header("depth", "0") + .match_header("authorization", auth_a.as_str()) + .with_status(207) + .with_header("content-type", "application/xml; charset=utf-8") + .with_body(principal("secondary@vandelay.org")) + .create(); + let _pb = server + .mock("PROPFIND", "/dav/cal/") + .match_header("depth", "0") + .match_header("authorization", auth_b.as_str()) + .with_status(207) + .with_header("content-type", "application/xml; charset=utf-8") + .with_body(principal("tertiary@vandelay.org")) + .create(); + + let archive: PathBuf = std::env::temp_dir().join(format!( + "vandelay-dav-srcchange-{}-{}.sqlite", + std::process::id(), + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_nanos() + )); + let _ = std::fs::remove_file(&archive); + + let common = |archive: &PathBuf| CommonConfig { + archive: archive.clone(), + threads: 1, + dry_run: false, + max_retries: 0, + allow_invalid_certs: true, + logger: Logger::from_flags(false, 0), + }; + let config = |user: &str, pass: &str| DavImportConfig { + kind: DavKindArg::Caldav, + url: format!("{url}/dav/cal/"), + auth: DavAuth::Basic { + user: user.to_owned(), + password: pass.to_owned(), + }, + allow_cleartext: true, + dav_connections: 1, + multiget_batch: 50, + allow_source_change: false, + }; + + run(common(&archive), config("secondary@vandelay.org", "passA")).expect("user A import ok"); + let err = run(common(&archive), config("tertiary@vandelay.org", "passB")).unwrap_err(); + let _ = std::fs::remove_file(&archive); + assert!( + matches!(err, vandelay::error::Error::SourceChange(_)), + "importing a different user into the same archive must trigger source-change protection; got {err:?}" + ); +} + #[test] fn source_change_protection_rejects_different_session_url_same_account() { use rusqlite::Connection; diff --git a/tests/mock_imap.rs b/tests/mock_imap.rs index 12b21ae..3617284 100644 --- a/tests/mock_imap.rs +++ b/tests/mock_imap.rs @@ -511,6 +511,43 @@ fn control_script_one_folder(uidvalidity: u32, uidnext: u32, uids: &'static [u32 }) } +fn control_script_present_flags( + uidvalidity: u32, + uidnext: u32, + uids: &'static [u32], + flags_reply: &'static [(u32, &'static str)], +) -> Script { + Box::new(move |conn: &mut MockConn| -> std::io::Result<()> { + auth_preamble(conn, "IMAP4rev2 LITERAL+ AUTH=PLAIN")?; + let (tag, cmd) = conn.read_command()?; + assert_eq!(cmd, "LIST \"\" \"*\""); + conn.write_line("* LIST () \"/\" \"INBOX\"")?; + conn.write_line(&format!("{tag} OK LIST done"))?; + let (tag, cmd) = conn.read_command()?; + assert_eq!(cmd, "LSUB \"\" \"*\""); + conn.write_line(&format!("{tag} OK LSUB done"))?; + let (tag, cmd) = conn.read_command()?; + assert_eq!(cmd, "SELECT \"INBOX\""); + write_select(conn, &tag, uidvalidity, uidnext, uids.len() as u32)?; + let (tag, cmd) = conn.read_command()?; + assert_eq!(cmd, "UID SEARCH ALL"); + let uid_strs: Vec = uids.iter().map(|u| u.to_string()).collect(); + conn.write_line(&format!("* SEARCH {}", uid_strs.join(" ")))?; + conn.write_line(&format!("{tag} OK SEARCH done"))?; + let (tag, cmd) = conn.read_command()?; + assert!( + cmd.starts_with("UID FETCH") && cmd.contains("(UID FLAGS)") && !cmd.contains("BODY"), + "expected body-less flags fetch on the present set, got {cmd}" + ); + for (uid, flags) in flags_reply { + conn.write_line(&format!("* {uid} FETCH (UID {uid} FLAGS ({flags}))"))?; + } + conn.write_line(&format!("{tag} OK FETCH done"))?; + drain_until_close(conn); + Ok(()) + }) +} + #[test] fn coordinator_imports_one_folder_one_message() { let server = MockImap::start_scripts(vec![ @@ -755,7 +792,7 @@ fn coordinator_present_run_is_convergent() { let mut scripts: Vec