Compare commits

...
11 Commits
Author SHA1 Message Date
jcoffey-dev 687027c7cb export: bring matched items up to date on every run
ci / test (pull_request) Skipped
github/ci (branch) GitHub Actions
ci / github (pull_request) Successful in 2m50s
ci / announce (pull_request) Skipped
Export matched each item against the target and then skipped it, so a second
run -- the usual final pass of a cutover -- never carried anything that had
changed at the source since the first: read and flagged state, moves between
folders, edited contacts, events and Sieve scripts. It reported them as
skipped and exited 0, while the usage guide said matched items were updated.

Matched items are now updated, with one batched /set per type:

- Email: keywords are set to the archive's, added and removed, compared
  case-insensitively. Memberships of folders this run migrated are added and
  removed to match; folders that exist only on the target are left alone,
  and a message is never left in no folder. Properties the server did not
  report are not touched.
- Contacts and events: when both copies carry `updated`, the archive's is
  written only if it is newer, compared as instants so an offset or a
  fraction of a second is not taken for a change; otherwise each property the
  archive writes is compared, and those that differ are sent whole.
- Sieve scripts: the target's copy is downloaded and compared byte for byte
  with what export would write -- after renaming Stalwart's vendor names for
  an inbuxa target -- and replaced when it differs, so a renamed script is
  not re-uploaded on every run.

Updated items are counted as `updated`; unchanged ones stay `skipped`. The
usage guide now describes this.
2026-09-30 11:37:39 -07:00
jcoffey-dev 14a797cd65 Merge pull request 'Time out every connection, and scope --allow-invalid-certs' (#6) from fix/connection-safety into main
ci / test (push) Skipped
github/ci (branch) GitHub Actions
ci / github (push) Successful in 3m31s
ci / announce (push) Skipped
2026-09-30 18:34:50 +00:00
jcoffey-dev 01ac1a5b3e Merge pull request 'export: write a message once, in every folder it was in' (#4) from fix/one-message-many-folders into main
ci / test (push) Skipped
github/ci (branch) GitHub Actions
ci / github (push) Canceled after 1m36s
ci / announce (push) Canceled after 0s
2026-09-30 18:33:09 +00:00
jcoffey-dev 234b3203d7 Scope --allow-invalid-certs to the server the user named
ci / test (pull_request) Skipped
github/ci (branch) GitHub Actions
ci / github (pull_request) Successful in 2m33s
ci / announce (pull_request) Skipped
The flag switched certificate checks off for every connection in the run.
That included the Microsoft sign-in endpoints, so a user passing it for a
self-signed source also sent refresh tokens, device codes and EWS client
secrets over unverified TLS. It also covered the export target, and any host
a server redirected to or named for its API, uploads or downloads.

It now applies only where the user pointed it: the host of --url, for the
source of an import or the target of an export. For an Exchange import with
no --url, it covers the mailbox's own domain, where on-premises Autodiscover
looks, and then only the EWS endpoint Autodiscover finds. The Microsoft and
Google sign-in and cloud hosts are always verified, with or without the flag.

Each HTTP client keeps a verifying agent and, only when the flag applies, a
second one that accepts invalid certificates, and picks per request by host.
The sign-in modules no longer take the flag at all. Autodiscover v2, which is
Microsoft's own service, is always verified.
2026-09-30 11:31:44 -07:00
jcoffey-dev 15b1cbd637 Merge pull request 'Create the archive readable by its owner only' (#5) from fix/archive-permissions into main
ci / test (push) Skipped
github/ci (branch) GitHub Actions
ci / github (push) Successful in 3m51s
ci / announce (push) Skipped
2026-09-30 18:29:04 +00:00
jcoffey-dev 67891994e6 Merge pull request 'Rename Stalwart's Sieve names for inbuxa, and fail when activation does' (#3) from fix/sieve-stalwart-names into main
ci / test (push) Skipped
github/ci (branch) GitHub Actions
ci / github (push) Canceled after 1m34s
ci / announce (push) Canceled after 0s
2026-09-30 18:27:28 +00:00
jcoffey-dev 59a8edbc7e Merge pull request 'docs: export matches and skips, it does not update' (#2) from docs/export-matching into main
ci / test (push) Canceled after 0s
ci / github (push) Canceled after 0s
ci / announce (push) Canceled after 0s
2026-09-30 18:27:26 +00:00
jcoffey-dev 38d32811e4 Create the archive readable by its owner only
ci / test (pull_request) Skipped
github/ci (branch) GitHub Actions
ci / github (pull_request) Successful in 2m32s
ci / announce (pull_request) Skipped
An archive holds a whole mailbox, its contacts and calendars, and it was
created with SQLite's default mode, readable by every user on the machine.
On Unix a new archive is now created with mode 0600 before SQLite opens it;
SQLite gives the -wal and -shm files the database file's mode, so they
follow, which a test confirms. An existing archive that others can read is
left as it is, with a warning naming it and the chmod that fixes it.
Nothing changes on Windows.
2026-09-30 11:26:18 -07:00
jcoffey-dev eb38603dbc Time out every HTTP connection instead of waiting forever
No ureq agent set a timeout, and ureq sets none by default, so a connection
dropped silently mid-transfer (a NAT or load-balancer idle drop) hung a JMAP,
DAV, EWS or Graph run, or an export, with no error, and the retry logic never
got a chance to run. IMAP and ManageSieve already had read timeouts.

Every agent now takes its settings from a new net module: 30s to connect,
60s to send the request headers, 5 minutes for the server's first byte, and
30 minutes for a whole response body, which ureq counts as one budget for the
body rather than per read: enough for the 512 MiB limit at about 300 KB/s.
A JMAP upload's send budget grows with its size, from a 2 minute floor at an
assumed 64 KiB/s worst case.

A timeout is a transport error, and every client already retries those, so a
stalled transfer is now abandoned and retried. Tests pin that down for each
client. The TLS setup the seven agents repeated moves to one helper.
2026-09-30 11:25:34 -07:00
jcoffey-dev 470fda6ac8 Rename Stalwart's Sieve names for inbuxa, and fail when activation does
ci / test (pull_request) Skipped
github/ci (branch) GitHub Actions
ci / github (pull_request) Successful in 2m11s
ci / announce (pull_request) Skipped
inbuxa accepts vnd.inbuxa.while and vnd.inbuxa.expressions, and names its
environment items vnd.inbuxa.*, with no alias for Stalwart's vnd.stalwart.*
names. A script carried over unchanged failed to compile on the target, and
the account was left with no filtering behind a single warning.

When the target lists vnd.inbuxa extensions in its sieveExtensions, export
now renames Stalwart's names in the three places they are names: the strings
of a require list, the item name of an environment test, and
${env.vnd.stalwart.*} references inside strings. A small tokenizer skips
comments and reads quoted strings and text: blocks whole, so other text that
happens to contain the name is copied unchanged. Each rename is printed.
Against a target without the inbuxa names nothing changes.

The activation call is now checked. A method error, or an active script
that could not be created, is logged as an error naming the script and
counted as a failure, so export exits non-zero instead of leaving
filtering off with one warning line.
2026-09-30 11:24:16 -07:00
jcoffey-dev 54fd18b410 docs: export matches and skips, it does not update
ci / test (pull_request) Skipped
github/ci (branch) GitHub Actions
ci / github (pull_request) Successful in 1m50s
ci / announce (pull_request) Skipped
usage.md said items that match are updated. Export skips a match: a
change made at the source after the first export -- read state, a folder
move, an edited event, contact or Sieve script -- does not reach the target
on a second one. Say so until updating on re-export is built.
2026-09-30 11:16:47 -07:00
30 changed files with 2264 additions and 333 deletions
+36 -9
View File
@@ -26,7 +26,26 @@ policy, TLS handling -- besides its own. The ones that matter most:
| `-v`, `-vv`, `-vvv` | Increase log verbosity. |
| `-q, --quiet` | Warnings and errors only. |
| `--max-retries <N>` | Max retries per request on transient failures (default 5). |
| `--allow-invalid-certs` | Accept self-signed / invalid TLS certs. |
| `--allow-invalid-certs` | Accept a self-signed or otherwise invalid certificate from the server named by `--url` (see below). |
### Invalid certificates
`--allow-invalid-certs` is for a server with a self-signed certificate, and
it applies to that server only: the host in `--url`, whether that is the
source of an import or the target of an export. For an Exchange import with
no `--url`, it covers the mailbox's own domain and the hosts under it, which
is where on-premises Autodiscover looks, and then only the EWS endpoint
Autodiscover finds.
Every other host is verified as usual, including a host the server
redirects to or names for its API, uploads or downloads. The Microsoft and
Google sign-in and cloud endpoints are always verified, with or without the
flag: a certificate that fails there is an attack or a broken network, never
a server to trust.
Connections time out rather than wait forever: 30 seconds to connect, 5
minutes for the server's first byte, and 30 minutes to read a whole
response. A timed-out request is retried like any other transient failure.
Secrets come from the `INBUXA_MIGRATE_*` environment variables or a prompt;
see [Credentials](../README.md#credentials). The command line takes them too,
@@ -192,15 +211,23 @@ inbuxa-migrate export \
Writes `ARCHIVE` into an account on a JMAP server, usually inbuxa. It keeps no
state of its own: every run matches the archive against the target afresh.
By default it only adds and updates -- items that match are updated, the
rest are created, and anything already on the target that the archive does
not cover is left alone.
By default it only adds and updates: what the target lacks is created, what
it has is brought up to date, and anything on the target that the archive
does not cover is left alone. So a second run after a later import carries
what changed at the source in between.
Email is matched by Message-ID, or without one by sender, subject, date and
recipients; where several messages share one, size decides. A message the
source kept in several folders -- IMAP and Maildir copies, Gmail labels -- is
written once, in all of them, and a later run adds any folder it is still
missing on the target.
- **Email** is matched by Message-ID, or without one by sender, subject, date
and recipients; where several messages share one, size decides. A message
the source kept in several folders -- IMAP and Maildir copies, Gmail labels
-- is written once, in all of them. On a match, its keywords (read,
flagged and the rest) are set to the archive's, and so are its memberships
of folders this run migrated. Folders that exist only on the target are
left alone, and a message is never left in no folder.
- **Contacts and events** are matched by UID. When both copies carry an
`updated` time, the archive's is written only if it is newer; otherwise
the properties that differ are written.
- **Sieve scripts** are matched by name, and the target's is replaced when
its content differs from the archive's.
`--prune` also deletes what is on the target and not in the archive. It asks
first; `--yes` answers for it, for scripts. Export speaks JMAP only.
+4 -1
View File
@@ -122,7 +122,10 @@ struct GlobalArgs {
)]
max_retries: u32,
#[arg(long, help = "Accept self-signed / invalid TLS certificates")]
#[arg(
long,
help = "Accept an invalid TLS certificate from the --url host only; sign-in endpoints are always verified"
)]
allow_invalid_certs: bool,
}
+51 -19
View File
@@ -14,7 +14,6 @@ use ureq::Agent;
use ureq::Body;
use ureq::config::{Config, RedirectAuthHeaders};
use ureq::http::{Method, Request, Response};
use ureq::tls::{RootCerts, TlsConfig};
use crate::dav::parse::{ControlStrippingReader, DavResponse, parse_multistatus};
use crate::dav::retry::{DavOutcome, classify};
@@ -22,6 +21,7 @@ use crate::jmap::error::JmapError;
use crate::jmap::http::{Auth, RetryPolicy, retry_after_header};
use crate::jmap::retry::{self, RateLimitState};
use crate::logging::{HttpCall, LEVEL_BODIES, LEVEL_DEFAULT, LEVEL_PROGRESS, Logger};
use crate::net::{CertOverride, tls, with_timeouts};
const MAX_BODY: u64 = 512 * 1024 * 1024;
const LONG_RETRY_THRESHOLD: Duration = Duration::from_secs(10);
@@ -48,6 +48,8 @@ pub struct MultiStatus {
struct Inner {
agent: Agent,
lax_agent: Option<Agent>,
certs: CertOverride,
auth: Auth,
retry: RetryPolicy,
rate_limit: RateLimitState,
@@ -57,31 +59,42 @@ struct Inner {
user_agent: String,
}
impl Inner {
/// The agent for `url`: the one that accepts invalid certificates only for
/// a host `--allow-invalid-certs` covers, and the verifying one otherwise.
fn agent_for(&self, url: &str) -> &Agent {
match &self.lax_agent {
Some(lax) if self.certs.allows(url) => lax,
_ => &self.agent,
}
}
}
#[derive(Clone)]
pub struct DavClient {
inner: Arc<Inner>,
}
impl DavClient {
pub fn new(auth: Auth, retry: RetryPolicy, allow_invalid_certs: bool) -> Self {
let config: Config = Config::builder()
.http_status_as_error(false)
.allow_non_standard_methods(true)
.max_redirects(0)
.redirect_auth_headers(RedirectAuthHeaders::SameHost)
.tls_config(
TlsConfig::builder()
.unversioned_rustls_crypto_provider(std::sync::Arc::new(
rustls::crypto::aws_lc_rs::default_provider(),
))
.root_certs(RootCerts::PlatformVerifier)
.disable_verification(allow_invalid_certs)
.build(),
pub fn new(auth: Auth, retry: RetryPolicy, certs: CertOverride) -> Self {
let build = |accept_invalid: bool| -> Agent {
let config: Config = with_timeouts!(
Config::builder()
.http_status_as_error(false)
.allow_non_standard_methods(true)
.max_redirects(0)
.redirect_auth_headers(RedirectAuthHeaders::SameHost)
.tls_config(tls(accept_invalid))
)
.build();
config.new_agent()
};
let lax_agent = certs.is_active().then(|| build(true));
DavClient {
inner: Arc::new(Inner {
agent: config.new_agent(),
agent: build(false),
lax_agent,
certs,
auth,
retry,
rate_limit: RateLimitState::new(),
@@ -822,7 +835,7 @@ impl DavClient {
let request = builder
.body(payload)
.map_err(|e| ureq::Error::Other(Box::new(std::io::Error::other(e))))?;
self.inner.agent.run(request)
self.inner.agent_for(req.url).run(request)
}
}
@@ -941,6 +954,25 @@ fn truncate(body: &[u8]) -> String {
#[cfg(test)]
mod tests {
use super::*;
use crate::net::CertOverride;
#[test]
fn every_timeout_is_a_retryable_transport_error() {
for t in [
ureq::Timeout::Connect,
ureq::Timeout::SendRequest,
ureq::Timeout::SendBody,
ureq::Timeout::RecvResponse,
ureq::Timeout::RecvBody,
] {
let err = map_ureq_error(ureq::Error::Timeout(t));
assert!(matches!(err, JmapError::Transport(_)), "{t:?} -> {err:?}");
assert!(
matches!(transport_disposition(&err), retry::Disposition::Retryable),
"{t:?} must be retried"
);
}
}
#[test]
fn client_constructs_cleanly() {
@@ -950,7 +982,7 @@ mod tests {
password: "p".into(),
},
RetryPolicy::new(3),
false,
CertOverride::none(),
);
assert_eq!(c.retries_observed(), 0);
assert_eq!(c.retry_after_sleeps(), 0);
@@ -963,7 +995,7 @@ mod tests {
token: "abc".into(),
},
RetryPolicy::new(0),
false,
CertOverride::none(),
);
let logger = c.logger();
assert_eq!(logger.level(), LEVEL_DEFAULT);
+95
View File
@@ -1,5 +1,6 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -9,6 +10,7 @@ use rusqlite::Connection;
pub const SCHEMA_SQL: &str = include_str!("schema.sql");
pub fn open(path: &std::path::Path) -> Result<Connection, OpenError> {
private_archive(path)?;
let conn = Connection::open(path)?;
apply_pragmas(&conn)?;
apply_schema(&conn)?;
@@ -72,6 +74,50 @@ fn ensure_graph_ids_accept_file_nodes(conn: &Connection) -> Result<(), OpenError
Ok(())
}
/// An archive holds a whole mailbox, so it is created readable by its owner
/// only. SQLite gives its `-wal` and `-shm` files the database file's mode,
/// so they follow. An existing archive others can read is left as it is,
/// with a warning and the command that fixes it.
#[cfg(unix)]
fn private_archive(path: &std::path::Path) -> Result<(), OpenError> {
use std::os::unix::fs::OpenOptionsExt;
match std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.mode(0o600)
.open(path)
{
Ok(_) => Ok(()),
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
if let Some(warning) = permission_warning(path) {
eprintln!("warning: {warning}");
}
Ok(())
}
Err(e) => Err(OpenError::Create(e)),
}
}
#[cfg(not(unix))]
fn private_archive(_path: &std::path::Path) -> Result<(), OpenError> {
Ok(())
}
/// The warning for an archive that someone other than its owner can read.
#[cfg(unix)]
fn permission_warning(path: &std::path::Path) -> Option<String> {
use std::os::unix::fs::PermissionsExt;
let mode = std::fs::metadata(path).ok()?.permissions().mode();
(mode & 0o077 != 0).then(|| {
format!(
"archive {} can be read by other users (mode {:o}); run: chmod 600 {}",
path.display(),
mode & 0o777,
path.display()
)
})
}
fn apply_pragmas(conn: &Connection) -> Result<(), OpenError> {
conn.pragma_update(None, "journal_mode", "WAL")?;
conn.pragma_update(None, "foreign_keys", "ON")?;
@@ -83,4 +129,53 @@ fn apply_pragmas(conn: &Connection) -> Result<(), OpenError> {
pub enum OpenError {
#[error("sqlite error: {0}")]
Sqlite(#[from] rusqlite::Error),
#[error("cannot create the archive: {0}")]
Create(std::io::Error),
}
#[cfg(all(test, unix))]
mod permission_tests {
use super::*;
use std::os::unix::fs::PermissionsExt;
fn mode(p: &std::path::Path) -> u32 {
std::fs::metadata(p).unwrap().permissions().mode() & 0o777
}
fn scratch(name: &str) -> std::path::PathBuf {
let dir =
std::env::temp_dir().join(format!("inbuxa-migrate-perm-{}-{name}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
dir.join("a.sqlite")
}
#[test]
fn a_new_archive_and_its_wal_are_private() {
let path = scratch("new");
let conn = open(&path).unwrap();
conn.execute_batch("CREATE TABLE t(x); INSERT INTO t VALUES (1);")
.unwrap();
assert_eq!(mode(&path), 0o600);
let wal = path.with_extension("sqlite-wal");
assert!(wal.exists(), "WAL mode writes a -wal file");
assert_eq!(mode(&wal), 0o600);
drop(conn);
let _ = std::fs::remove_dir_all(path.parent().unwrap());
}
#[test]
fn an_existing_readable_archive_is_warned_about_not_changed() {
let path = scratch("existing");
drop(open(&path).unwrap());
std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o644)).unwrap();
let warning = permission_warning(&path).expect("warns");
assert!(warning.contains("chmod 600"), "{warning}");
assert!(warning.contains("644"), "{warning}");
drop(open(&path).unwrap());
assert_eq!(mode(&path), 0o644, "left as it is");
std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o600)).unwrap();
assert!(permission_warning(&path).is_none());
let _ = std::fs::remove_dir_all(path.parent().unwrap());
}
}
+23 -19
View File
@@ -1,5 +1,6 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -9,10 +10,10 @@ use quick_xml::events::Event;
use serde_json::Value;
use ureq::Agent;
use ureq::config::Config;
use ureq::tls::{RootCerts, TlsConfig};
use crate::exchange_ews::error::EwsError;
use crate::exchange_ews::parse::entity_to_char;
use crate::net::{CertOverride, tls, with_timeouts};
const V2_HOST: &str = "https://outlook.office365.com";
const POX_REQ_NS: &str =
@@ -48,7 +49,7 @@ pub fn discover(
supplied_url: Option<&str>,
email: Option<&str>,
auth_header: Option<&str>,
allow_invalid_certs: bool,
certs: &CertOverride,
) -> Result<DiscoveryResult, EwsError> {
if let Some(url) = supplied_url
&& is_fully_qualified_ews_url(url)
@@ -63,8 +64,17 @@ pub fn discover(
"either a fully-qualified --url or --mailbox is required".to_owned(),
));
};
let agent = build_agent(allow_invalid_certs);
if let Ok(url) = autodiscover_v2(&agent, email) {
// Autodiscover v2 is Microsoft's own service and is always verified; a v1
// candidate gets the relaxed agent only if `--allow-invalid-certs` covers it.
let strict = build_agent(false);
let lax = certs.is_active().then(|| build_agent(true));
let pick = |url: &str| -> &Agent {
match &lax {
Some(agent) if certs.allows(url) => agent,
_ => &strict,
}
};
if let Ok(url) = autodiscover_v2(&strict, email) {
return Ok(DiscoveryResult {
ews_url: url,
source: DiscoverySource::V2,
@@ -82,7 +92,7 @@ pub fn discover(
let candidates = pox_candidates(domain);
for candidate in &candidates {
tried.push(candidate.clone());
match autodiscover_v1(&agent, candidate, &current_email, auth_header) {
match autodiscover_v1(pick(candidate), candidate, &current_email, auth_header) {
Ok(PoxOutcome::EwsUrl(url)) => {
return Ok(DiscoveryResult {
ews_url: url,
@@ -104,7 +114,7 @@ pub fn discover(
url_redirects += 1;
tried.push(url.clone());
if let Ok(PoxOutcome::EwsUrl(u)) =
autodiscover_v1(&agent, &url, &current_email, auth_header)
autodiscover_v1(pick(&url), &url, &current_email, auth_header)
{
return Ok(DiscoveryResult {
ews_url: u,
@@ -130,19 +140,13 @@ pub fn discover(
)))
}
fn build_agent(allow_invalid_certs: bool) -> Agent {
let config: Config = Config::builder()
.http_status_as_error(false)
.tls_config(
TlsConfig::builder()
.unversioned_rustls_crypto_provider(std::sync::Arc::new(
rustls::crypto::aws_lc_rs::default_provider(),
))
.root_certs(RootCerts::PlatformVerifier)
.disable_verification(allow_invalid_certs)
.build(),
)
.build();
fn build_agent(accept_invalid: bool) -> Agent {
let config: Config = with_timeouts!(
Config::builder()
.http_status_as_error(false)
.tls_config(tls(accept_invalid))
)
.build();
config.new_agent()
}
+47 -18
View File
@@ -12,7 +12,6 @@ use std::time::{Duration, Instant};
use ureq::Agent;
use ureq::config::{Config, RedirectAuthHeaders};
use ureq::tls::{RootCerts, TlsConfig};
use crate::exchange_ews::error::EwsError;
use crate::exchange_ews::parse::{EnvelopeKind, SoapFault, read_envelope_summary};
@@ -22,12 +21,15 @@ use crate::exchange_ews::types::ServerVersion;
use crate::jmap::http::{Auth, RetryPolicy, retry_after_header};
use crate::jmap::retry::{self, Disposition, RateLimitState};
use crate::logging::{HttpCall, LEVEL_BODIES, LEVEL_DEFAULT, LEVEL_PROGRESS, Logger};
use crate::net::{CertOverride, tls, with_timeouts};
const MAX_BODY: u64 = 2 * 1024 * 1024 * 1024;
const LONG_RETRY_THRESHOLD: Duration = Duration::from_secs(10);
struct Inner {
agent: Agent,
lax_agent: Option<Agent>,
certs: CertOverride,
auth: Mutex<Auth>,
impersonated_smtp: Mutex<Option<String>>,
anchor_mailbox: Mutex<Option<String>>,
@@ -42,6 +44,17 @@ struct Inner {
user_agent: String,
}
impl Inner {
/// The agent for `url`: the one that accepts invalid certificates only for
/// a host `--allow-invalid-certs` covers, and the verifying one otherwise.
fn agent_for(&self, url: &str) -> &Agent {
match &self.lax_agent {
Some(lax) if self.certs.allows(url) => lax,
_ => &self.agent,
}
}
}
#[derive(Clone)]
pub struct EwsClient {
inner: Arc<Inner>,
@@ -54,23 +67,23 @@ pub struct SoapResponse {
}
impl EwsClient {
pub fn new(auth: Auth, retry: RetryPolicy, allow_invalid_certs: bool) -> EwsClient {
let config: Config = Config::builder()
.http_status_as_error(false)
.redirect_auth_headers(RedirectAuthHeaders::SameHost)
.tls_config(
TlsConfig::builder()
.unversioned_rustls_crypto_provider(std::sync::Arc::new(
rustls::crypto::aws_lc_rs::default_provider(),
))
.root_certs(RootCerts::PlatformVerifier)
.disable_verification(allow_invalid_certs)
.build(),
pub fn new(auth: Auth, retry: RetryPolicy, certs: CertOverride) -> EwsClient {
let build = |accept_invalid: bool| -> Agent {
let config: Config = with_timeouts!(
Config::builder()
.http_status_as_error(false)
.redirect_auth_headers(RedirectAuthHeaders::SameHost)
.tls_config(tls(accept_invalid))
)
.build();
config.new_agent()
};
let lax_agent = certs.is_active().then(|| build(true));
EwsClient {
inner: Arc::new(Inner {
agent: config.new_agent(),
agent: build(false),
lax_agent,
certs,
auth: Mutex::new(auth),
impersonated_smtp: Mutex::new(None),
anchor_mailbox: Mutex::new(None),
@@ -408,7 +421,7 @@ impl EwsClient {
fn one_attempt(&self, url: &str, body: &str, action: &str) -> AttemptOutcome {
let mut req = self
.inner
.agent
.agent_for(url)
.post(url)
.header("Authorization", self.auth_header())
.header("Content-Type", "text/xml; charset=utf-8")
@@ -535,13 +548,29 @@ fn truncate(body: &[u8]) -> String {
#[cfg(test)]
mod tests {
use super::*;
use crate::net::CertOverride;
#[test]
fn every_timeout_is_a_transport_error_and_so_retried() {
// Every EwsError::Transport goes round the retry loop in `execute`.
for t in [
ureq::Timeout::Connect,
ureq::Timeout::SendRequest,
ureq::Timeout::SendBody,
ureq::Timeout::RecvResponse,
ureq::Timeout::RecvBody,
] {
let err = map_ureq_error(ureq::Error::Timeout(t));
assert!(matches!(err, EwsError::Transport(_)), "{t:?} -> {err:?}");
}
}
#[test]
fn client_constructs_with_defaults() {
let c = EwsClient::new(
Auth::Bearer { token: "t".into() },
RetryPolicy::new(3),
false,
CertOverride::none(),
);
assert_eq!(c.server_version(), ServerVersion::Exchange2013Sp1);
assert_eq!(c.retries_observed(), 0);
@@ -553,7 +582,7 @@ mod tests {
let c = EwsClient::new(
Auth::Bearer { token: "t".into() },
RetryPolicy::new(0),
false,
CertOverride::none(),
);
c.set_server_version(ServerVersion::Exchange2019);
assert_eq!(c.server_version(), ServerVersion::Exchange2019);
@@ -564,7 +593,7 @@ mod tests {
let c = EwsClient::new(
Auth::Bearer { token: "t".into() },
RetryPolicy::new(0),
false,
CertOverride::none(),
);
c.set_anchor_mailbox(Some("alice@x".to_owned()));
assert_eq!(c.anchor_header().as_deref(), Some("alice@x"));
+19 -35
View File
@@ -1,5 +1,6 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -10,9 +11,9 @@ use std::time::Duration;
use encodify::base64::{Base64, Padding, URL_SAFE};
use serde_json::Value;
use ureq::config::Config;
use ureq::tls::{RootCerts, TlsConfig};
use crate::exchange_ews::error::EwsError;
use crate::net::{tls, with_timeouts};
pub const SCOPE_APP_ONLY: &str = "https://outlook.office365.com/.default";
pub const SCOPE_DELEGATED: &str =
@@ -46,7 +47,7 @@ pub enum OAuthFlow {
},
}
pub fn acquire(flow: &OAuthFlow, allow_invalid_certs: bool) -> Result<AcquiredToken, EwsError> {
pub fn acquire(flow: &OAuthFlow) -> Result<AcquiredToken, EwsError> {
match flow {
OAuthFlow::PreAcquired { token } => {
let claims = decode_jwt_claims(token).unwrap_or_default();
@@ -63,10 +64,8 @@ pub fn acquire(flow: &OAuthFlow, allow_invalid_certs: bool) -> Result<AcquiredTo
tenant,
client_id,
client_secret,
} => client_credentials(tenant, client_id, client_secret, allow_invalid_certs),
OAuthFlow::DeviceCode { tenant, client_id } => {
device_code_flow(tenant, client_id, allow_invalid_certs)
}
} => client_credentials(tenant, client_id, client_secret),
OAuthFlow::DeviceCode { tenant, client_id } => device_code_flow(tenant, client_id),
}
}
@@ -104,19 +103,13 @@ fn device_code_endpoint(tenant: &str) -> String {
format!("https://login.microsoftonline.com/{tenant}/oauth2/v2.0/devicecode")
}
fn build_agent(allow_invalid_certs: bool) -> ureq::Agent {
let config: Config = Config::builder()
.http_status_as_error(false)
.tls_config(
TlsConfig::builder()
.unversioned_rustls_crypto_provider(std::sync::Arc::new(
rustls::crypto::aws_lc_rs::default_provider(),
))
.root_certs(RootCerts::PlatformVerifier)
.disable_verification(allow_invalid_certs)
.build(),
)
.build();
fn build_agent() -> ureq::Agent {
let config: Config = with_timeouts!(
Config::builder()
.http_status_as_error(false)
.tls_config(tls(false))
)
.build();
config.new_agent()
}
@@ -124,9 +117,8 @@ fn client_credentials(
tenant: &str,
client_id: &str,
client_secret: &str,
allow_invalid_certs: bool,
) -> Result<AcquiredToken, EwsError> {
let agent = build_agent(allow_invalid_certs);
let agent = build_agent();
let body = form_encode(&[
("client_id", client_id),
("client_secret", client_secret),
@@ -142,12 +134,8 @@ fn client_credentials(
parse_token_response(resp)
}
fn device_code_flow(
tenant: &str,
client_id: &str,
allow_invalid_certs: bool,
) -> Result<AcquiredToken, EwsError> {
let agent = build_agent(allow_invalid_certs);
fn device_code_flow(tenant: &str, client_id: &str) -> Result<AcquiredToken, EwsError> {
let agent = build_agent();
let body = form_encode(&[("client_id", client_id), ("scope", SCOPE_DELEGATED)]);
let endpoint = device_code_endpoint(tenant);
let mut resp = agent
@@ -279,9 +267,8 @@ pub fn refresh_with_token(
tenant: &str,
client_id: &str,
refresh_token: &str,
allow_invalid_certs: bool,
) -> Result<AcquiredToken, EwsError> {
let agent = build_agent(allow_invalid_certs);
let agent = build_agent();
let body = form_encode(&[
("client_id", client_id),
("grant_type", "refresh_token"),
@@ -366,12 +353,9 @@ mod tests {
#[test]
fn pre_acquired_flow_decodes_claims() {
let token = make_jwt("t-2", "bob@x", 9999999999);
let acq = acquire(
&OAuthFlow::PreAcquired {
token: token.clone(),
},
false,
)
let acq = acquire(&OAuthFlow::PreAcquired {
token: token.clone(),
})
.unwrap();
assert_eq!(acq.access_token, token);
assert_eq!(acq.tenant_id.as_deref(), Some("t-2"));
+50 -17
View File
@@ -13,7 +13,6 @@ use std::time::{Duration, Instant};
use serde_json::Value;
use ureq::Agent;
use ureq::config::{Config, RedirectAuthHeaders};
use ureq::tls::{RootCerts, TlsConfig};
use ureq::{ResponseExt, http::Uri};
use crate::exchange_graph::error::GraphError;
@@ -21,6 +20,7 @@ use crate::exchange_graph::retry::{HttpClass, classify_http_status, is_throttled
use crate::jmap::http::{RetryPolicy, cross_host, retry_after_header};
use crate::jmap::retry::{self, RateLimitState};
use crate::logging::{HttpCall, LEVEL_BODIES, LEVEL_DEFAULT, LEVEL_PROGRESS, Logger};
use crate::net::{CertOverride, tls, with_timeouts};
const MAX_BODY: u64 = 256 * 1024 * 1024;
const LONG_RETRY_THRESHOLD: Duration = Duration::from_secs(10);
@@ -63,6 +63,8 @@ impl GraphResponse {
struct Inner {
agent: Agent,
lax_agent: Option<Agent>,
certs: CertOverride,
bearer: Mutex<String>,
retry: RetryPolicy,
rate_limit: RateLimitState,
@@ -73,6 +75,17 @@ struct Inner {
user_agent: String,
}
impl Inner {
/// The agent for `url`: the one that accepts invalid certificates only for
/// a host `--allow-invalid-certs` covers, and the verifying one otherwise.
fn agent_for(&self, url: &str) -> &Agent {
match &self.lax_agent {
Some(lax) if self.certs.allows(url) => lax,
_ => &self.agent,
}
}
}
#[derive(Clone)]
pub struct GraphClient {
inner: Arc<Inner>,
@@ -89,23 +102,23 @@ enum Attempt {
}
impl GraphClient {
pub fn new(bearer: String, retry: RetryPolicy, allow_invalid_certs: bool) -> GraphClient {
let config: Config = Config::builder()
.http_status_as_error(false)
.redirect_auth_headers(RedirectAuthHeaders::SameHost)
.tls_config(
TlsConfig::builder()
.unversioned_rustls_crypto_provider(std::sync::Arc::new(
rustls::crypto::aws_lc_rs::default_provider(),
))
.root_certs(RootCerts::PlatformVerifier)
.disable_verification(allow_invalid_certs)
.build(),
pub fn new(bearer: String, retry: RetryPolicy, certs: CertOverride) -> GraphClient {
let build = |accept_invalid: bool| -> Agent {
let config: Config = with_timeouts!(
Config::builder()
.http_status_as_error(false)
.redirect_auth_headers(RedirectAuthHeaders::SameHost)
.tls_config(tls(accept_invalid))
)
.build();
config.new_agent()
};
let lax_agent = certs.is_active().then(|| build(true));
GraphClient {
inner: Arc::new(Inner {
agent: config.new_agent(),
agent: build(false),
lax_agent,
certs,
bearer: Mutex::new(bearer),
retry,
rate_limit: RateLimitState::new(),
@@ -310,7 +323,7 @@ impl GraphClient {
extra_prefer: &[&str],
) -> Attempt {
let mut req = match method {
"GET" => self.inner.agent.get(url),
"GET" => self.inner.agent_for(url).get(url),
other => {
return Attempt::Transport(GraphError::Connect(format!(
"unsupported method {other} (graph importer is read-only)"
@@ -474,10 +487,30 @@ fn format_retry_wait(d: Duration) -> String {
#[cfg(test)]
mod tests {
use super::*;
use crate::net::CertOverride;
#[test]
fn every_timeout_is_a_transport_error_and_so_retried() {
// `execute` retries every GraphError::Transport; only Connect is fatal.
for t in [
ureq::Timeout::Connect,
ureq::Timeout::SendRequest,
ureq::Timeout::SendBody,
ureq::Timeout::RecvResponse,
ureq::Timeout::RecvBody,
] {
let err = map_ureq_error(ureq::Error::Timeout(t));
assert!(matches!(err, GraphError::Transport(_)), "{t:?} -> {err:?}");
}
}
#[test]
fn defaults_construct() {
let c = GraphClient::new("token".to_owned(), RetryPolicy::new(3), false);
let c = GraphClient::new(
"token".to_owned(),
RetryPolicy::new(3),
CertOverride::none(),
);
assert_eq!(c.retries_observed(), 0);
assert_eq!(c.retry_after_sleeps(), 0);
assert_eq!(c.requests_observed(), 0);
@@ -486,7 +519,7 @@ mod tests {
#[test]
fn bearer_can_be_swapped_at_runtime() {
let c = GraphClient::new("old".to_owned(), RetryPolicy::new(0), false);
let c = GraphClient::new("old".to_owned(), RetryPolicy::new(0), CertOverride::none());
c.set_bearer("new".to_owned());
assert_eq!(c.auth_header(), "Bearer new");
}
+14 -24
View File
@@ -1,5 +1,6 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -10,9 +11,9 @@ use std::time::{Duration, Instant};
use encodify::base64::{Base64, Padding, URL_SAFE};
use serde_json::Value;
use ureq::config::Config;
use ureq::tls::{RootCerts, TlsConfig};
use crate::exchange_graph::error::GraphError;
use crate::net::{tls, with_timeouts};
pub const SCOPES: &str =
"offline_access User.Read Mail.Read MailboxSettings.Read Calendars.Read Contacts.Read";
@@ -78,29 +79,23 @@ pub struct AcquiredToken {
pub name: Option<String>,
}
fn build_agent(allow_invalid_certs: bool) -> ureq::Agent {
let config: Config = Config::builder()
.http_status_as_error(false)
.tls_config(
TlsConfig::builder()
.unversioned_rustls_crypto_provider(std::sync::Arc::new(
rustls::crypto::aws_lc_rs::default_provider(),
))
.root_certs(RootCerts::PlatformVerifier)
.disable_verification(allow_invalid_certs)
.build(),
)
.build();
fn build_agent() -> ureq::Agent {
let config: Config = with_timeouts!(
Config::builder()
.http_status_as_error(false)
.tls_config(tls(false))
)
.build();
config.new_agent()
}
pub fn acquire(flow: &OAuthFlow, allow_invalid_certs: bool) -> Result<AcquiredToken, GraphError> {
pub fn acquire(flow: &OAuthFlow) -> Result<AcquiredToken, GraphError> {
match flow {
OAuthFlow::PreAcquired { token } => Ok(token_from_string(token.clone())),
OAuthFlow::DeviceCode {
authority,
client_id,
} => device_code_flow(authority, client_id, allow_invalid_certs),
} => device_code_flow(authority, client_id),
}
}
@@ -213,12 +208,8 @@ pub fn parse_token_response(status: u16, json: &Value) -> TokenResponse {
}
}
fn device_code_flow(
authority: &str,
client_id: &str,
allow_invalid_certs: bool,
) -> Result<AcquiredToken, GraphError> {
let agent = build_agent(allow_invalid_certs);
fn device_code_flow(authority: &str, client_id: &str) -> Result<AcquiredToken, GraphError> {
let agent = build_agent();
let body = form_encode(&[("client_id", client_id), ("scope", SCOPES)]);
let endpoint = device_code_endpoint(authority);
let mut resp = agent
@@ -294,9 +285,8 @@ pub fn refresh_access_token(
authority: &str,
client_id: &str,
refresh_token: &str,
allow_invalid_certs: bool,
) -> Result<AcquiredToken, GraphError> {
let agent = build_agent(allow_invalid_certs);
let agent = build_agent();
let body = form_encode(&[
("client_id", client_id),
("grant_type", "refresh_token"),
+3 -1
View File
@@ -1,5 +1,6 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -170,6 +171,7 @@ fn extract_account_id(principal: &Value, name: &str) -> Result<String, Error> {
#[cfg(test)]
mod tests {
use super::*;
use crate::net::CertOverride;
fn session_with(name: &str, id: &str) -> Session {
let raw = serde_json::json!({
@@ -186,7 +188,7 @@ mod tests {
HttpClient::new(
crate::jmap::http::Auth::Bearer { token: "t".into() },
crate::jmap::http::RetryPolicy::new(0),
false,
CertOverride::none(),
)
}
+102 -29
View File
@@ -1,5 +1,6 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -13,7 +14,6 @@ use encodify::base64::STANDARD;
use serde_json::Value;
use ureq::Agent;
use ureq::config::{Config, RedirectAuthHeaders};
use ureq::tls::{RootCerts, TlsConfig};
use ureq::{ResponseExt, http::Uri};
use crate::jmap::error::JmapError;
@@ -21,6 +21,7 @@ use crate::jmap::inflight::{Permit, Semaphore};
use crate::jmap::retry::{self, Disposition, RateLimitState};
use crate::jmap::session::Limits;
use crate::logging::{HttpCall, LEVEL_BODIES, LEVEL_DEFAULT, LEVEL_PROGRESS, Logger};
use crate::net::{CertOverride, send_body_budget, tls, with_timeouts};
const MAX_BODY: u64 = 512 * 1024 * 1024;
@@ -69,9 +70,10 @@ impl RetryPolicy {
struct Inner {
agent: Agent,
lax_agent: Option<Agent>,
certs: CertOverride,
auth: Auth,
retry: RetryPolicy,
allow_invalid_certs: bool,
rate_limit: RateLimitState,
log_level: AtomicU8,
requests_gate: OnceLock<Semaphore>,
@@ -81,6 +83,17 @@ struct Inner {
retry_after_sleeps: AtomicU64,
}
impl Inner {
/// The agent for `url`: the one that accepts invalid certificates only for
/// a host `--allow-invalid-certs` covers, and the verifying one otherwise.
fn agent_for(&self, url: &str) -> &Agent {
match &self.lax_agent {
Some(lax) if self.certs.allows(url) => lax,
_ => &self.agent,
}
}
}
#[derive(Debug, Clone, Copy)]
enum Kind {
Api,
@@ -103,26 +116,25 @@ enum Attempt {
}
impl HttpClient {
pub fn new(auth: Auth, retry: RetryPolicy, allow_invalid_certs: bool) -> Self {
let config: Config = Config::builder()
.http_status_as_error(false)
.redirect_auth_headers(RedirectAuthHeaders::SameHost)
.tls_config(
TlsConfig::builder()
.unversioned_rustls_crypto_provider(std::sync::Arc::new(
rustls::crypto::aws_lc_rs::default_provider(),
))
.root_certs(RootCerts::PlatformVerifier)
.disable_verification(allow_invalid_certs)
.build(),
pub fn new(auth: Auth, retry: RetryPolicy, certs: CertOverride) -> Self {
let build = |accept_invalid: bool| -> Agent {
let config: Config = with_timeouts!(
Config::builder()
.http_status_as_error(false)
.redirect_auth_headers(RedirectAuthHeaders::SameHost)
.tls_config(tls(accept_invalid))
)
.build();
config.new_agent()
};
let lax_agent = certs.is_active().then(|| build(true));
HttpClient {
inner: Arc::new(Inner {
agent: config.new_agent(),
agent: build(false),
lax_agent,
certs,
auth,
retry,
allow_invalid_certs,
rate_limit: RateLimitState::new(),
log_level: AtomicU8::new(LEVEL_DEFAULT),
requests_gate: OnceLock::new(),
@@ -172,10 +184,6 @@ impl HttpClient {
&self.inner.retry
}
pub fn allow_invalid_certs(&self) -> bool {
self.inner.allow_invalid_certs
}
pub fn rate_limit(&self) -> &RateLimitState {
&self.inner.rate_limit
}
@@ -386,17 +394,22 @@ impl HttpClient {
let result = if let Some(payload) = body {
let mut req = self
.inner
.agent
.agent_for(url)
.post(url)
.header("Authorization", auth)
.header("Accept", "application/json");
if let Some(ct) = content_type {
req = req.header("Content-Type", ct);
}
req.send(payload)
// A blob upload can run to hundreds of megabytes, so its send
// budget grows with its size instead of the agent's flat default.
req.config()
.timeout_send_body(Some(send_body_budget(payload.len())))
.build()
.send(payload)
} else {
self.inner
.agent
.agent_for(url)
.get(url)
.header("Authorization", auth)
.header("Accept", "application/json")
@@ -644,6 +657,66 @@ pub fn format_retry_wait(d: Duration) -> String {
#[cfg(test)]
mod tests {
use super::*;
use crate::net::CertOverride;
#[test]
fn invalid_certificates_are_accepted_only_for_the_named_host() {
let client = HttpClient::new(
Auth::Bearer {
token: "t".to_owned(),
},
RetryPolicy::new(0),
CertOverride::for_url(true, "https://mail.example.test/.well-known/jmap"),
);
let inner = &client.inner;
let lax = inner.lax_agent.as_ref().expect("a relaxed agent exists");
assert!(std::ptr::eq(
inner.agent_for("https://mail.example.test/api"),
lax
));
assert!(std::ptr::eq(
inner.agent_for("https://files.example.test/upload"),
&inner.agent
));
assert!(std::ptr::eq(
inner.agent_for("https://login.microsoftonline.com/common/oauth2/v2.0/token"),
&inner.agent
));
}
#[test]
fn without_the_flag_there_is_no_relaxed_agent() {
let client = HttpClient::new(
Auth::Bearer {
token: "t".to_owned(),
},
RetryPolicy::new(0),
CertOverride::for_url(false, "https://mail.example.test/"),
);
assert!(client.inner.lax_agent.is_none());
assert!(std::ptr::eq(
client.inner.agent_for("https://mail.example.test/api"),
&client.inner.agent
));
}
#[test]
fn every_timeout_is_a_retryable_transport_error() {
for t in [
ureq::Timeout::Connect,
ureq::Timeout::SendRequest,
ureq::Timeout::SendBody,
ureq::Timeout::RecvResponse,
ureq::Timeout::RecvBody,
] {
let err = map_ureq_error(ureq::Error::Timeout(t));
assert!(matches!(err, JmapError::Transport(_)), "{t:?} -> {err:?}");
assert!(
matches!(transport_disposition(&err), Disposition::Retryable),
"{t:?} must be retried"
);
}
}
#[test]
fn basic_header_matches_rfc7617_example() {
@@ -711,7 +784,7 @@ mod tests {
token: "t".to_owned(),
},
RetryPolicy::new(0),
false,
CertOverride::none(),
);
let body = br#"{"type":"urn:ietf:params:jmap:error:limit","limit":"someServerLimit"}"#;
assert!(matches!(
@@ -727,7 +800,7 @@ mod tests {
token: "t".to_owned(),
},
RetryPolicy::new(0),
false,
CertOverride::none(),
);
let body = br#"{"type":"urn:ietf:params:jmap:error:limit","limit":"maxSizeRequest"}"#;
assert!(matches!(
@@ -743,7 +816,7 @@ mod tests {
token: "t".to_owned(),
},
RetryPolicy::new(0),
false,
CertOverride::none(),
);
let body =
br#"{"type":"urn:ietf:params:jmap:error:limit","limit":"maxConcurrentRequests"}"#;
@@ -772,7 +845,7 @@ mod tests {
token: "t".to_owned(),
},
RetryPolicy::new(0),
false,
CertOverride::none(),
);
client.set_limits(&limits_with(10, 4, 4));
let err = client
@@ -796,7 +869,7 @@ mod tests {
token: "t".to_owned(),
},
RetryPolicy::new(0),
false,
CertOverride::none(),
);
client.set_limits(&limits_with(1024, 4, 4));
let err = client
@@ -848,7 +921,7 @@ mod tests {
token: "t".to_owned(),
},
RetryPolicy::new(0),
false,
CertOverride::none(),
);
assert_eq!(client.retries_observed(), 0);
assert_eq!(client.retry_after_sleeps(), 0);
+2 -1
View File
@@ -786,6 +786,7 @@ fn decode_set(mr: &MethodCall) -> SetOutcome {
#[cfg(test)]
mod tests {
use crate::net::CertOverride;
use std::cell::Cell;
use super::*;
@@ -798,7 +799,7 @@ mod tests {
token: "t".to_owned(),
},
RetryPolicy::new(max_retries),
false,
CertOverride::none(),
)
}
+2
View File
@@ -1,5 +1,6 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -16,6 +17,7 @@ pub mod inspect;
pub mod jmap;
pub mod logging;
pub mod managesieve;
pub mod net;
pub mod secret;
pub mod sync;
pub mod types;
+284
View File
@@ -0,0 +1,284 @@
/*
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
//! Settings every HTTP agent shares: timeouts, TLS, and which hosts, if any,
//! may present a certificate that does not verify.
use std::time::Duration;
use ureq::tls::{RootCerts, TlsConfig};
/// Opening the socket and completing any TLS handshake.
pub const CONNECT: Duration = Duration::from_secs(30);
/// Writing the request line and headers.
pub const SEND_REQUEST: Duration = Duration::from_secs(60);
/// Waiting for the response headers once the request is sent. This is the
/// server's thinking time: a large `Email/import`, an EWS `FindItem` over a big
/// folder or a CalDAV REPORT can legitimately take a while before the first
/// byte comes back.
pub const RECV_RESPONSE: Duration = Duration::from_secs(5 * 60);
/// Reading the whole response body. ureq counts this as one budget for the
/// entire body, not per read, so it has to cover the largest body a client
/// accepts (512 MiB) on a slow link: 30 minutes is about 300 KB/s. A stalled
/// transfer is abandoned and retried after at most this long.
pub const RECV_BODY: Duration = Duration::from_secs(30 * 60);
/// Sending a request body when its size is not known in advance. Uploads know
/// their size and get [`send_body_budget`] instead.
pub const SEND_BODY: Duration = Duration::from_secs(30 * 60);
/// The slowest upload rate a send budget allows for, in bytes per second.
const MIN_UPLOAD_RATE: u64 = 64 * 1024;
/// The floor under every send budget, so small bodies still get a sensible
/// allowance on a slow or busy connection.
const SEND_BODY_FLOOR: Duration = Duration::from_secs(2 * 60);
/// How long sending a body of `len` bytes may take: the floor plus the time it
/// takes at [`MIN_UPLOAD_RATE`].
pub fn send_body_budget(len: usize) -> Duration {
SEND_BODY_FLOOR + Duration::from_secs(len as u64 / MIN_UPLOAD_RATE)
}
/// Applies the shared timeouts to a ureq `ConfigBuilder`. A macro rather than
/// a function because ureq keeps the builder's scope types private, so a
/// function could not name them.
macro_rules! with_timeouts {
($builder:expr) => {
$builder
.timeout_connect(Some($crate::net::CONNECT))
.timeout_send_request(Some($crate::net::SEND_REQUEST))
.timeout_send_body(Some($crate::net::SEND_BODY))
.timeout_recv_response(Some($crate::net::RECV_RESPONSE))
.timeout_recv_body(Some($crate::net::RECV_BODY))
};
}
pub(crate) use with_timeouts;
/// TLS settings for an agent: the platform's roots, and certificate checks off
/// only when `accept_invalid` is set.
pub fn tls(accept_invalid: bool) -> TlsConfig {
TlsConfig::builder()
.unversioned_rustls_crypto_provider(std::sync::Arc::new(
rustls::crypto::aws_lc_rs::default_provider(),
))
.root_certs(RootCerts::PlatformVerifier)
.disable_verification(accept_invalid)
.build()
}
/// Hosts that are always verified, whatever `--allow-invalid-certs` says:
/// the Microsoft and Google sign-in and cloud endpoints. A certificate that
/// fails there is an attack or a broken network, never a self-signed server
/// the user meant to trust. Matched as a suffix on a label boundary.
const ALWAYS_VERIFY: &[&str] = &[
"microsoftonline.com",
"microsoftonline.us",
"microsoft.com",
"microsoft.us",
"office365.com",
"office.com",
"outlook.com",
"chinacloudapi.cn",
"partner.outlook.cn",
"google.com",
"googleapis.com",
"gmail.com",
];
/// Where `--allow-invalid-certs` applies: the host the user named, or, for
/// Exchange Autodiscover without a `--url`, the mailbox's own domain and its
/// subdomains. Everything else, including any host a server redirects or
/// points to, is verified as usual.
#[derive(Debug, Clone, Default)]
pub struct CertOverride {
hosts: Vec<String>,
domains: Vec<String>,
}
impl CertOverride {
/// Verify everything.
pub fn none() -> Self {
Self::default()
}
/// When `enabled`, accept invalid certificates from the host of `url`.
pub fn for_url(enabled: bool, url: &str) -> Self {
match (enabled, host_of(url)) {
(true, Some(host)) if !always_verified(&host) => CertOverride {
hosts: vec![host],
domains: Vec::new(),
},
_ => Self::none(),
}
}
/// When `enabled`, accept invalid certificates from `domain` and every
/// host under it.
pub fn for_domain(enabled: bool, domain: &str) -> Self {
let domain = domain.trim_end_matches('.').to_ascii_lowercase();
if enabled && !domain.is_empty() && !always_verified(&domain) {
CertOverride {
hosts: Vec::new(),
domains: vec![domain],
}
} else {
Self::none()
}
}
/// The same override, narrowed to the host of `url`, if `url` is one this
/// override already covers. Used once Autodiscover has found the real
/// endpoint.
pub fn narrowed_to(&self, url: &str) -> Self {
match host_of(url) {
Some(host) if self.allows_host(&host) => CertOverride {
hosts: vec![host],
domains: Vec::new(),
},
_ => Self::none(),
}
}
/// Whether this override covers anything at all.
pub fn is_active(&self) -> bool {
!self.hosts.is_empty() || !self.domains.is_empty()
}
/// Whether a certificate that does not verify is accepted for `url`.
pub fn allows(&self, url: &str) -> bool {
host_of(url).is_some_and(|host| self.allows_host(&host))
}
fn allows_host(&self, host: &str) -> bool {
if always_verified(host) {
return false;
}
self.hosts.iter().any(|h| h == host)
|| self
.domains
.iter()
.any(|d| host == d || host.ends_with(&format!(".{d}")))
}
}
fn host_of(url: &str) -> Option<String> {
let parsed = url::Url::parse(url).ok()?;
let host = parsed
.host_str()?
.trim_end_matches('.')
.to_ascii_lowercase();
Some(
host.trim_start_matches('[')
.trim_end_matches(']')
.to_owned(),
)
}
fn always_verified(host: &str) -> bool {
ALWAYS_VERIFY
.iter()
.any(|d| host == *d || host.ends_with(&format!(".{d}")))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn disabled_flag_covers_nothing() {
let o = CertOverride::for_url(false, "https://mail.example.test/jmap");
assert!(!o.is_active());
assert!(!o.allows("https://mail.example.test/jmap"));
}
#[test]
fn covers_only_the_named_host() {
let o = CertOverride::for_url(true, "https://Mail.Example.test:8443/.well-known/jmap");
assert!(o.is_active());
assert!(o.allows("https://mail.example.test/api"));
assert!(o.allows("https://MAIL.example.test:9000/upload"));
assert!(!o.allows("https://files.example.test/download"));
assert!(!o.allows("https://example.test/"));
assert!(!o.allows("https://mail.example.test.evil.test/"));
}
#[test]
fn sign_in_and_cloud_hosts_are_always_verified() {
for url in [
"https://login.microsoftonline.com/common/oauth2/v2.0/token",
"https://graph.microsoft.com/v1.0/me",
"https://outlook.office365.com/EWS/Exchange.asmx",
"https://autodiscover-s.outlook.com/autodiscover/autodiscover.xml",
"https://oauth2.googleapis.com/token",
"https://accounts.google.com/o/oauth2/device/code",
] {
let o = CertOverride::for_url(true, url);
assert!(!o.is_active(), "{url}");
assert!(!o.allows(url), "{url}");
}
}
#[test]
fn a_domain_covers_its_subdomains_but_not_look_alikes() {
let o = CertOverride::for_domain(true, "Corp.Example.");
assert!(o.allows("https://autodiscover.corp.example/autodiscover/autodiscover.xml"));
assert!(o.allows("https://corp.example/autodiscover/autodiscover.xml"));
assert!(!o.allows("https://notcorp.example/"));
assert!(!o.allows("https://corp.example.evil.test/"));
}
#[test]
fn a_domain_override_never_reaches_microsoft() {
let o = CertOverride::for_domain(true, "office365.com");
assert!(!o.is_active());
let corp = CertOverride::for_domain(true, "corp.example");
assert!(!corp.allows("https://outlook.office365.com/EWS/Exchange.asmx"));
}
#[test]
fn narrowing_keeps_only_a_covered_endpoint() {
let o = CertOverride::for_domain(true, "corp.example");
let inside = o.narrowed_to("https://mail.corp.example/EWS/Exchange.asmx");
assert!(inside.allows("https://mail.corp.example/EWS/Exchange.asmx"));
assert!(!inside.allows("https://autodiscover.corp.example/"));
let outside = o.narrowed_to("https://outlook.office365.com/EWS/Exchange.asmx");
assert!(!outside.is_active());
}
#[test]
fn ip_literals_are_matched() {
let o = CertOverride::for_url(true, "https://[::1]:8443/jmap");
assert!(o.allows("https://[::1]:9000/other"));
let v4 = CertOverride::for_url(true, "https://192.0.2.10/jmap");
assert!(v4.allows("https://192.0.2.10:8443/"));
assert!(!v4.allows("https://192.0.2.11/"));
}
#[test]
fn send_budget_grows_with_size() {
assert_eq!(send_body_budget(0), Duration::from_secs(120));
assert_eq!(
send_body_budget(64 * 1024 * 600),
Duration::from_secs(120 + 600)
);
assert!(send_body_budget(512 * 1024 * 1024) > Duration::from_secs(2 * 60 * 60));
}
#[test]
fn timeouts_are_applied_to_a_config() {
let config: ureq::config::Config = with_timeouts!(ureq::config::Config::builder()).build();
let t = config.timeouts();
assert_eq!(t.connect, Some(CONNECT));
assert_eq!(t.send_request, Some(SEND_REQUEST));
assert_eq!(t.send_body, Some(SEND_BODY));
assert_eq!(t.recv_response, Some(RECV_RESPONSE));
assert_eq!(t.recv_body, Some(RECV_BODY));
}
}
+90 -1
View File
@@ -1,10 +1,11 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
use std::collections::HashMap;
use std::collections::{HashMap, HashSet};
use std::io::{IsTerminal, Write};
use rusqlite::Connection;
@@ -48,6 +49,14 @@ impl Maps {
fn insert(&mut self, ty: ObjectType, local: i64, target: JmapId) {
self.m.entry(ty).or_default().insert(local, target);
}
/// Every target id this run mapped for `ty`: the objects it migrated.
fn targets_of(&self, ty: ObjectType) -> HashSet<String> {
self.m
.get(&ty)
.map(|m| m.values().map(|id| id.0.clone()).collect())
.unwrap_or_default()
}
}
impl TargetResolver for Maps {
@@ -97,6 +106,34 @@ impl<'a> Uploader<'a> {
Ok(id)
}
/// As `upload_with`, but sends `bytes` in place of the stored blob: for
/// content rewritten on its way to the target. Cached under the same
/// local id, so a retry sends the rewritten bytes again.
fn upload_bytes_as(
&mut self,
local_id: i64,
content_type: &str,
bytes: &[u8],
) -> Result<JmapId, JmapError> {
self.touched.push(local_id);
if let Some(id) = self.cache.get(&local_id) {
return Ok(id.clone());
}
let id = if self.net.dry_run {
JmapId(format!("dryrun-blob-{local_id}"))
} else {
blobxfer::upload_bytes(
&self.net.client,
&self.net.session,
&self.net.account,
content_type,
bytes,
)?
};
self.cache.insert(local_id, id.clone());
Ok(id)
}
fn invalidate(&mut self, local_id: i64) {
self.cache.remove(&local_id);
}
@@ -435,6 +472,8 @@ mod keyed;
mod sieve;
mod sieve_names;
mod uidtype;
mod email;
@@ -502,6 +541,56 @@ mod common {
)
}
/// Sends `updates` (target id, patch) as batched `/set` calls and counts
/// the result into `counts`. A dry run counts them as updated and sends
/// nothing.
pub fn update_batch(
net: &Net,
ty: ObjectType,
updates: Vec<(String, Value)>,
counts: &mut TypeCounts,
logger: &Logger,
) {
if updates.is_empty() {
return;
}
if net.dry_run {
counts.updated += updates.len() as u64;
return;
}
let total = updates.len() as u64;
let mut map = Map::new();
for (id, patch) in updates {
map.insert(id, patch);
}
match set_call(
&net.client,
&net.api,
&net.account,
ty.jmap_name(),
SetRequest {
update: Some(Value::Object(map)),
..Default::default()
},
&net.limits,
) {
Ok(outcome) => {
counts.updated += outcome.updated.len() as u64;
for (id, err) in &outcome.not_updated {
logger.warn(&format!("{}/set {id} not updated: {err}", ty.jmap_name()));
counts.failed += 1;
}
}
Err(e) => {
logger.warn(&format!(
"{}/set: updating {total} object(s) failed: {e}",
ty.jmap_name()
));
counts.failed += total;
}
}
}
fn blob_not_found(outcome: &crate::jmap::request::SetOutcome, cid: &str) -> bool {
outcome.not_created.iter().any(|(c, err)| {
c == cid && err.get("type").and_then(Value::as_str) == Some("blobNotFound")
+141 -67
View File
@@ -9,12 +9,12 @@ use std::collections::{HashMap, HashSet};
use serde_json::{Map, Value, json};
use super::common::{jid, target_query_get};
use super::common::{jid, target_query_get, update_batch};
use super::{Maps, Net, Plan, Uploader};
use crate::error::Error;
use crate::jmap::error::JmapError;
use crate::jmap::request::{
MethodCall, Request, SetRequest, check_method_error, get_objects, retry_method_call, set_call,
MethodCall, Request, check_method_error, get_objects, retry_method_call,
};
use crate::jmap::retry::MethodCallKind;
use crate::jmap::wire::JmapId;
@@ -62,8 +62,12 @@ pub fn reconcile(
) -> Result<Plan, Error> {
let ty = ObjectType::Email;
let target_min = target_query_get(net, ty, Some(&["messageId", "size", "mailboxIds"]))
.map_err(Error::from)?;
let target_min = target_query_get(
net,
ty,
Some(&["messageId", "size", "mailboxIds", "keywords"]),
)
.map_err(Error::from)?;
let mut indices: Vec<EmailIndex> = target_min.iter().map(server_index).collect();
let fallback_ids: Vec<JmapId> = target_min
@@ -128,11 +132,12 @@ pub fn reconcile(
.collect();
let pairs = pair_with_targets(&local_keys, &sizes, &target_keys, &targets);
let mut membership_updates: Vec<(String, Value)> = Vec::new();
let migrated = maps.targets_of(ObjectType::Mailbox);
let mut updates: Vec<(String, Value)> = Vec::new();
for (i, unit) in units.iter().enumerate() {
match pairs[i] {
Some(t) => match missing_memberships(&unit.row, &targets[t], maps) {
Some(patch) => membership_updates.push((targets[t].id.clone(), patch)),
Some(t) => match email_patch(&unit.row, &targets[t], maps, &migrated) {
Some(patch) => updates.push((targets[t].id.clone(), patch)),
None => counts.skipped += 1,
},
None => export_one(
@@ -146,7 +151,7 @@ pub fn reconcile(
),
}
}
send_membership_updates(net, membership_updates, counts, logger);
update_batch(net, ty, updates, counts, logger);
Ok(Plan::default())
}
@@ -197,6 +202,7 @@ struct TargetEmail {
id: String,
size: Option<u64>,
mailboxes: Option<HashSet<String>>,
keywords: Option<HashSet<String>>,
}
impl TargetEmail {
@@ -208,6 +214,10 @@ impl TargetEmail {
.get("mailboxIds")
.and_then(Value::as_object)
.map(|m| m.keys().cloned().collect()),
keywords: v
.get("keywords")
.and_then(Value::as_object)
.map(|m| m.keys().map(|k| k.to_lowercase()).collect()),
}
}
}
@@ -254,65 +264,58 @@ fn pair_with_targets(
out
}
/// The `Email/set` patch adding the folders a matched email is missing on the
/// target, or `None` when it is already in all of them. Folders that exist
/// only on the target are left alone.
fn missing_memberships(row: &EmailRow, target: &TargetEmail, maps: &Maps) -> Option<Value> {
let have = target.mailboxes.as_ref()?;
/// The `Email/set` patch that brings a matched email on the target in line
/// with the archive, or `None` when it already is. The source is taken as
/// the truth for what it covers: keywords are added and removed to match, and
/// so are memberships of folders this run migrated. Folders that exist only
/// on the target are left alone, an email is never left in no folder, and
/// whatever the server did not report is not touched.
fn email_patch(
row: &EmailRow,
target: &TargetEmail,
maps: &Maps,
migrated: &HashSet<String>,
) -> Option<Value> {
let mut patch = Map::new();
for ml in &row.mailbox_locals {
if let Some(t) = maps.target(ObjectType::Mailbox, *ml)
&& !have.contains(&t.0)
{
patch.insert(format!("mailboxIds/{}", t.0), Value::Bool(true));
if let Some(have) = &target.mailboxes {
let want: HashSet<String> = row
.mailbox_locals
.iter()
.filter_map(|ml| maps.target(ObjectType::Mailbox, *ml).map(|t| t.0))
.collect();
let add: Vec<&String> = want.iter().filter(|t| !have.contains(*t)).collect();
let remove: Vec<&String> = have
.iter()
.filter(|t| migrated.contains(*t) && !want.contains(*t))
.collect();
let left = have.len() - remove.len() + add.len();
for t in add {
patch.insert(
format!("mailboxIds/{}", pointer_escape(t)),
Value::Bool(true),
);
}
if left > 0 {
for t in remove {
patch.insert(format!("mailboxIds/{}", pointer_escape(t)), Value::Null);
}
}
}
if let Some(have) = &target.keywords {
let want: HashSet<String> = row.keywords.iter().map(|k| k.to_lowercase()).collect();
for k in want.difference(have) {
patch.insert(format!("keywords/{}", pointer_escape(k)), Value::Bool(true));
}
for k in have.difference(&want) {
patch.insert(format!("keywords/{}", pointer_escape(k)), Value::Null);
}
}
(!patch.is_empty()).then_some(Value::Object(patch))
}
fn send_membership_updates(
net: &Net,
updates: Vec<(String, Value)>,
counts: &mut TypeCounts,
logger: &Logger,
) {
if updates.is_empty() {
return;
}
if net.dry_run {
counts.updated += updates.len() as u64;
return;
}
let total = updates.len() as u64;
let mut map = Map::new();
for (id, patch) in updates {
map.insert(id, patch);
}
match set_call(
&net.client,
&net.api,
&net.account,
ObjectType::Email.jmap_name(),
SetRequest {
update: Some(Value::Object(map)),
..Default::default()
},
&net.limits,
) {
Ok(outcome) => {
counts.updated += outcome.updated.len() as u64;
for (id, err) in &outcome.not_updated {
logger.warn(&format!("Email/set {id}: folders not added: {err}"));
counts.failed += 1;
}
}
Err(e) => {
logger.warn(&format!(
"Email/set: adding folders to {total} email(s) failed: {e}"
));
counts.failed += total;
}
}
/// Escapes one JSON Pointer segment (RFC 6901), as JMAP patch paths use.
fn pointer_escape(segment: &str) -> String {
segment.replace('~', "~0").replace('/', "~1")
}
fn build_mailbox_ids(row: &EmailRow, maps: &Maps) -> Option<Map<String, Value>> {
@@ -570,9 +573,14 @@ mod tests {
id: id.to_owned(),
size,
mailboxes: mailboxes.map(|m| m.iter().map(|s| (*s).to_owned()).collect()),
keywords: None,
}
}
fn set(items: &[&str]) -> HashSet<String> {
items.iter().map(|s| (*s).to_owned()).collect()
}
fn mid(m: &str) -> EmailKey {
EmailKey::MessageId(m.to_owned())
}
@@ -628,19 +636,85 @@ mod tests {
assert_eq!(pairs, vec![Some(0), None]);
}
#[test]
fn missing_memberships_adds_only_the_absent_migrated_folders() {
fn two_folder_maps() -> Maps {
let mut maps = Maps::default();
maps.insert(ObjectType::Mailbox, 1, JmapId("T1".into()));
maps.insert(ObjectType::Mailbox, 2, JmapId("T2".into()));
maps
}
#[test]
fn patch_adds_the_absent_migrated_folders() {
let maps = two_folder_maps();
let r = row(10, &[1, 2, 3], &[]);
let patch = missing_memberships(&r, &target("E", None, Some(&["T1", "Own"])), &maps)
.expect("T2 is missing");
let patch = email_patch(
&r,
&target("E", None, Some(&["T1", "Own"])),
&maps,
&set(&["T1", "T2"]),
)
.expect("T2 is missing");
assert_eq!(patch, json!({ "mailboxIds/T2": true }));
assert!(missing_memberships(&r, &target("E", None, Some(&["T1", "T2"])), &maps).is_none());
assert!(
missing_memberships(&r, &target("E", None, None), &maps).is_none(),
email_patch(
&r,
&target("E", None, Some(&["T1", "T2"])),
&maps,
&set(&["T1", "T2"])
)
.is_none()
);
assert!(
email_patch(&r, &target("E", None, None), &maps, &set(&["T1", "T2"])).is_none(),
"unknown membership is left alone"
);
}
#[test]
fn patch_moves_between_migrated_folders_but_keeps_target_only_ones() {
let maps = two_folder_maps();
let r = row(10, &[2], &[]);
let patch = email_patch(
&r,
&target("E", None, Some(&["T1", "Own"])),
&maps,
&set(&["T1", "T2"]),
)
.unwrap();
assert_eq!(
patch,
json!({ "mailboxIds/T2": true, "mailboxIds/T1": null })
);
}
#[test]
fn patch_never_leaves_an_email_in_no_folder() {
let maps = two_folder_maps();
let r = row(10, &[9], &[]);
assert!(
email_patch(
&r,
&target("E", None, Some(&["T1"])),
&maps,
&set(&["T1", "T2"])
)
.is_none(),
"the only folder is not removed when nothing replaces it"
);
}
#[test]
fn patch_syncs_keywords_both_ways_case_insensitively() {
let maps = two_folder_maps();
let r = row(10, &[1], &["$Seen", "work/urgent"]);
let mut t = target("E", None, Some(&["T1"]));
t.keywords = Some(set(&["$seen", "$flagged"]));
let patch = email_patch(&r, &t, &maps, &set(&["T1", "T2"])).unwrap();
assert_eq!(
patch,
json!({ "keywords/work~1urgent": true, "keywords/$flagged": null })
);
t.keywords = Some(set(&["$seen", "work/urgent"]));
assert!(email_patch(&r, &t, &maps, &set(&["T1", "T2"])).is_none());
}
}
+144 -10
View File
@@ -1,5 +1,6 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -8,11 +9,14 @@ use std::collections::{HashMap, HashSet};
use serde_json::{Value, json};
use super::common::{create_batch, jid, retry_if_blob_missing, target_get_all};
use super::common::{create_batch, jid, retry_if_blob_missing, target_get_all, update_batch};
use super::sieve_names;
use super::{Maps, Net, Plan, Uploader};
use crate::error::Error;
use crate::jmap::request::Request;
use crate::logging::Logger;
use crate::jmap::blobxfer;
use crate::jmap::request::{Request, check_method_error};
use crate::logging::{LEVEL_DEFAULT, Logger};
use crate::sync::import_jmap::mapping::BlobBytes;
use crate::sync::import_jmap::mapping::{SIEVE_SELECT, row_to_sieve_script};
use crate::sync::{Context, TypeCounts};
use crate::types::ObjectType;
@@ -28,10 +32,14 @@ pub fn reconcile(
let targets = target_get_all(net, ty).map_err(Error::from)?;
let mut target_by_name: HashMap<String, String> = HashMap::new();
let mut target_blob: HashMap<String, String> = HashMap::new();
for t in &targets {
let (Some(id), Some(name)) = (jid(t), t.get("name").and_then(Value::as_str)) else {
continue;
};
if let Some(blob) = t.get("blobId").and_then(Value::as_str) {
target_blob.insert(id.clone(), blob.to_owned());
}
target_by_name.insert(name.to_owned(), id);
}
@@ -57,18 +65,67 @@ pub fn reconcile(
let mut active_target: Option<String> = None;
let mut deactivate = false;
let mut uploader = Uploader::new(net, &ctx.conn);
let rename_vendor = sieve_names::target_uses_inbuxa_names(&target_sieve_extensions(net));
let wanted_active = locals
.iter()
.find(|(_, _, a, _)| *a)
.map(|(_, n, _, _)| n.clone().unwrap_or_default());
let mut updates: Vec<(String, Value)> = Vec::new();
for (local, name, is_active, blob_local) in &locals {
let matched = name.as_ref().and_then(|n| target_by_name.get(n)).cloned();
let label = name.as_deref().unwrap_or("(unnamed)");
let rewritten = if rename_vendor {
renamed_script(&uploader, *blob_local)?
} else {
None
};
let target_id = if let Some(id) = matched {
counts.skipped += 1;
// Compare what would be written -- the renamed bytes where the
// script needed renaming -- so an unchanged script stays unchanged.
let ours = match &rewritten {
Some((bytes, _)) => bytes.clone(),
None => uploader.bytes(*blob_local).map_err(Error::from)?,
};
match content_differs(net, &ours, target_blob.get(&id)) {
Ok(false) => counts.skipped += 1,
Ok(true) => {
let blob = match &rewritten {
Some((bytes, renamed)) => {
log_renames(label, renamed, logger);
uploader.upload_bytes_as(*blob_local, "application/sieve", bytes)
}
None => uploader.upload_with(*blob_local, "application/sieve"),
};
match blob {
Ok(b) => updates.push((id.clone(), json!({ "blobId": b.0 }))),
Err(e) => {
logger.warn(&format!(
"SieveScript {label}: upload for update failed: {e}"
));
counts.failed += 1;
}
}
}
Err(e) => {
logger.warn(&format!("SieveScript {label}: not compared: {e}"));
counts.skipped += 1;
}
}
id
} else {
let cid = format!("c{local}");
if let Some((_, renamed)) = &rewritten {
log_renames(label, renamed, logger);
}
let rewritten = rewritten.as_ref().map(|(bytes, _)| bytes);
let build = |up: &mut Uploader<'_>| -> Result<Value, Error> {
let blob_id = up
.upload_with(*blob_local, "application/sieve")
.map_err(Error::from)?;
let blob_id = match &rewritten {
Some(bytes) => up.upload_bytes_as(*blob_local, "application/sieve", bytes),
None => up.upload_with(*blob_local, "application/sieve"),
}
.map_err(Error::from)?;
let mut obj = serde_json::Map::new();
if let Some(n) = name {
obj.insert("name".to_owned(), Value::String(n.clone()));
@@ -92,7 +149,7 @@ pub fn reconcile(
}
None => {
for (cid, err) in &outcome.not_created {
logger.warn(&format!("SieveScript {cid} not created: {err}"));
logger.warn(&format!("SieveScript {label} ({cid}) not created: {err}"));
}
counts.failed += 1;
continue;
@@ -104,9 +161,20 @@ pub fn reconcile(
}
}
update_batch(net, ty, updates, counts, logger);
if active_target.is_none() && locals.iter().all(|(_, _, a, _)| !*a) {
deactivate = true;
}
if let (Some(name), None) = (&wanted_active, &active_target) {
// The script that was active at the source never made it to the
// target (its creation failure is already counted): say plainly that
// the account now has no filtering, rather than leave it to a warning.
logger.error(&format!(
"the active Sieve script \"{name}\" could not be created on the target; \
no filtering is active there"
));
}
if !net.dry_run {
let mut req = Request::new();
@@ -118,8 +186,20 @@ pub fn reconcile(
json!({ "accountId": net.account })
};
req.call("SieveScript/set", args, "a");
if let Err(e) = req.send(&net.client, &net.api) {
logger.warn(&format!("SieveScript activation failed: {e}"));
let result = req
.send(&net.client, &net.api)
.and_then(|resp| resp.by_call_id("a").cloned())
.and_then(|mr| check_method_error(&mr));
if let Err(e) = result {
match (&active_target, &wanted_active) {
(Some(_), Some(name)) => {
logger.error(&format!(
"the Sieve script \"{name}\" was created but could not be activated: {e}"
));
counts.failed += 1;
}
_ => logger.warn(&format!("SieveScript activation failed: {e}")),
}
}
}
@@ -136,3 +216,57 @@ pub fn reconcile(
active_sieve_target: active_target,
})
}
/// The target's `sieveExtensions`, from its Sieve account capability.
fn target_sieve_extensions(net: &Net) -> Vec<String> {
net.session
.account_capabilities(&net.account)
.and_then(|caps| caps.get("urn:ietf:params:jmap:sieve"))
.and_then(|c| c.get("sieveExtensions"))
.and_then(Value::as_array)
.map(|a| {
a.iter()
.filter_map(Value::as_str)
.map(str::to_owned)
.collect()
})
.unwrap_or_default()
}
/// A script's bytes after renaming, and the names that were renamed.
type Renamed = (Vec<u8>, Vec<String>);
/// The script's bytes with Stalwart's vendor names renamed for an inbuxa
/// target, and the names renamed, or `None` when it needs no change.
fn renamed_script(uploader: &Uploader<'_>, blob_local: i64) -> Result<Option<Renamed>, Error> {
let bytes = uploader.bytes(blob_local).map_err(Error::from)?;
Ok(sieve_names::rewrite(&bytes))
}
/// Prints each rename made to a script about to be written.
fn log_renames(label: &str, renamed: &[String], logger: &Logger) {
for old in renamed {
let new = old.replacen("vnd.stalwart.", "vnd.inbuxa.", 1);
if logger.enabled(LEVEL_DEFAULT) {
eprintln!("export: SieveScript {label}: renamed {old} to {new}");
}
}
}
/// Whether the target's copy of a script differs from `ours`. A target that
/// reports no blob is taken as different, so ours is written.
fn content_differs(net: &Net, ours: &[u8], target_blob: Option<&String>) -> Result<bool, Error> {
let Some(target_blob) = target_blob else {
return Ok(true);
};
let theirs = blobxfer::download_bytes(
&net.client,
&net.session,
&net.account,
target_blob,
"application/sieve",
"script.sieve",
)
.map_err(Error::from)?;
Ok(ours != theirs.as_slice())
}
+342
View File
@@ -0,0 +1,342 @@
/*
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
//! Stalwart's vendor Sieve names, renamed for inbuxa.
//!
//! inbuxa accepts `vnd.inbuxa.while` and `vnd.inbuxa.expressions` where
//! Stalwart accepted `vnd.stalwart.*`, with no alias, and names its
//! environment items the same way. A script carried over unchanged fails to
//! compile on inbuxa, so export renames those names -- and only those names:
//! the strings of a `require` list, the name argument of an `environment`
//! test, and `${env.vnd.stalwart.…}` references inside strings. Everything
//! else in the script, including other strings that happen to contain the
//! text, is copied byte for byte.
const OLD: &str = "vnd.stalwart.";
const NEW: &str = "vnd.inbuxa.";
const OLD_ENV_REF: &str = "${env.vnd.stalwart.";
const NEW_ENV_REF: &str = "${env.vnd.inbuxa.";
/// Whether the target advertises inbuxa's vendor extensions, from the
/// `sieveExtensions` list of its `urn:ietf:params:jmap:sieve` account
/// capability.
pub fn target_uses_inbuxa_names(sieve_extensions: &[String]) -> bool {
sieve_extensions.iter().any(|e| e.starts_with(NEW))
}
/// The script with Stalwart's vendor names renamed, and the old names that
/// were changed, in order. `None` when nothing needed renaming, or when the
/// script is not UTF-8 (left alone rather than guessed at).
pub fn rewrite(script: &[u8]) -> Option<(Vec<u8>, Vec<String>)> {
let text = std::str::from_utf8(script).ok()?;
if !text.contains(OLD) {
return None;
}
let mut out = String::with_capacity(text.len());
let mut renamed = Vec::new();
let mut copied = 0;
let mut context = Context::None;
for tok in Tokens::new(text) {
match tok.kind {
Kind::Word => {
context = match tok.text(text).to_ascii_lowercase().as_str() {
"require" => Context::Require,
"environment" => Context::Environment,
_ if context == Context::Environment => Context::Environment,
_ => Context::None,
};
}
Kind::Tag => {
// `environment :comparator "i;octet"`: the comparator's own
// string is not the item name.
if context == Context::Environment
&& tok.text(text).eq_ignore_ascii_case(":comparator")
{
context = Context::EnvironmentComparator;
}
}
Kind::Quoted | Kind::Multiline => {
let (start, end) = tok.content;
let content = &text[start..end];
let mut replacement: Option<String> = None;
let whole_name = matches!(context, Context::Require | Context::Environment);
if whole_name && content.starts_with(OLD) {
replacement = Some(format!("{NEW}{}", &content[OLD.len()..]));
renamed.push(content.to_owned());
}
let current = replacement.as_deref().unwrap_or(content);
if current.contains(OLD_ENV_REF) {
let mut n = 0;
let mut rest = current;
while let Some(i) = rest.find(OLD_ENV_REF) {
let tail = &rest[i + 2..];
let name_end = tail.find('}').unwrap_or(tail.len());
renamed.push(tail[4..name_end].to_owned());
rest = &rest[i + OLD_ENV_REF.len()..];
n += 1;
}
if n > 0 {
replacement = Some(current.replace(OLD_ENV_REF, NEW_ENV_REF));
}
}
if let Some(r) = replacement {
out.push_str(&text[copied..start]);
out.push_str(&r);
copied = end;
}
context = match context {
Context::Require => Context::Require,
Context::EnvironmentComparator => Context::Environment,
_ => Context::None,
};
}
Kind::Punct(';') | Kind::Punct('{') | Kind::Punct('}') => context = Context::None,
Kind::Punct(_) => {}
}
}
if renamed.is_empty() {
return None;
}
out.push_str(&text[copied..]);
Some((out.into_bytes(), renamed))
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum Context {
None,
Require,
Environment,
EnvironmentComparator,
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum Kind {
Word,
Tag,
Quoted,
Multiline,
Punct(char),
}
struct Token {
kind: Kind,
span: (usize, usize),
/// The string's content, without quotes or the `text:` framing.
content: (usize, usize),
}
impl Token {
fn text<'a>(&self, src: &'a str) -> &'a str {
&src[self.span.0..self.span.1]
}
}
/// Just enough of RFC 5228's lexer to find strings and the words before
/// them: comments are skipped, and quoted strings and `text:` blocks are
/// read whole, so nothing inside them is mistaken for a command.
struct Tokens<'a> {
src: &'a str,
pos: usize,
}
impl<'a> Tokens<'a> {
fn new(src: &'a str) -> Self {
Tokens { src, pos: 0 }
}
}
impl Iterator for Tokens<'_> {
type Item = Token;
fn next(&mut self) -> Option<Token> {
let b = self.src.as_bytes();
loop {
while self.pos < b.len() && b[self.pos].is_ascii_whitespace() {
self.pos += 1;
}
if self.pos >= b.len() {
return None;
}
if b[self.pos] == b'#' {
while self.pos < b.len() && b[self.pos] != b'\n' {
self.pos += 1;
}
continue;
}
if b[self.pos..].starts_with(b"/*") {
self.pos = match self.src[self.pos + 2..].find("*/") {
Some(i) => self.pos + 2 + i + 2,
None => b.len(),
};
continue;
}
break;
}
let start = self.pos;
let c = b[start];
if c == b'"' {
let mut i = start + 1;
while i < b.len() && b[i] != b'"' {
i += if b[i] == b'\\' { 2 } else { 1 };
}
let end = i.min(b.len());
self.pos = (end + 1).min(b.len());
return Some(Token {
kind: Kind::Quoted,
span: (start, self.pos),
content: (start + 1, end),
});
}
if c.is_ascii_alphabetic() || c == b'_' || c == b':' {
let mut i = start + 1;
while i < b.len() && (b[i].is_ascii_alphanumeric() || b[i] == b'_') {
i += 1;
}
let word = &self.src[start..i];
if word.eq_ignore_ascii_case("text:")
|| (word.eq_ignore_ascii_case("text") && b.get(i) == Some(&b':'))
{
let after = if b.get(i) == Some(&b':') { i + 1 } else { i };
// The body starts after the rest of the `text:` line and runs
// to a line holding a single dot.
let body = match self.src[after..].find('\n') {
Some(n) => after + n + 1,
None => b.len(),
};
let (body_end, next) = find_dot_line(self.src, body);
self.pos = next;
return Some(Token {
kind: Kind::Multiline,
span: (start, next),
content: (body, body_end),
});
}
self.pos = i;
let kind = if c == b':' { Kind::Tag } else { Kind::Word };
return Some(Token {
kind,
span: (start, i),
content: (start, i),
});
}
let ch = self.src[start..].chars().next().unwrap_or('\0');
self.pos = start + ch.len_utf8();
Some(Token {
kind: Kind::Punct(ch),
span: (start, self.pos),
content: (start, self.pos),
})
}
}
/// End of a `text:` body (the start of its closing dot line) and the offset
/// after that line.
fn find_dot_line(src: &str, from: usize) -> (usize, usize) {
let mut line_start = from;
while line_start < src.len() {
let line_end = src[line_start..]
.find('\n')
.map(|n| line_start + n)
.unwrap_or(src.len());
if src[line_start..line_end].trim_end_matches('\r') == "." {
return (line_start, (line_end + 1).min(src.len()));
}
line_start = line_end + 1;
}
(src.len(), src.len())
}
#[cfg(test)]
mod tests {
use super::*;
fn run(s: &str) -> Option<(String, Vec<String>)> {
rewrite(s.as_bytes()).map(|(b, r)| (String::from_utf8(b).unwrap(), r))
}
#[test]
fn renames_a_require_list() {
let (out, renamed) =
run("require [\"fileinto\", \"vnd.stalwart.while\", \"vnd.stalwart.expressions\"];\n")
.unwrap();
assert_eq!(
out,
"require [\"fileinto\", \"vnd.inbuxa.while\", \"vnd.inbuxa.expressions\"];\n"
);
assert_eq!(renamed, ["vnd.stalwart.while", "vnd.stalwart.expressions"]);
}
#[test]
fn renames_a_single_require_string() {
let (out, _) = run("REQUIRE \"vnd.stalwart.while\";").unwrap();
assert_eq!(out, "REQUIRE \"vnd.inbuxa.while\";");
}
#[test]
fn renames_the_environment_item_name_only() {
let src = "if environment :comparator \"i;octet\" :is \"vnd.stalwart.username\" \"vnd.stalwart.x\" { keep; }";
let (out, renamed) = run(src).unwrap();
assert_eq!(
out,
"if environment :comparator \"i;octet\" :is \"vnd.inbuxa.username\" \"vnd.stalwart.x\" { keep; }"
);
assert_eq!(renamed, ["vnd.stalwart.username"]);
}
#[test]
fn renames_env_references_inside_strings() {
let src = "set \"box\" \"${env.vnd.stalwart.default_mailbox}/Archive\";";
let (out, renamed) = run(src).unwrap();
assert_eq!(
out,
"set \"box\" \"${env.vnd.inbuxa.default_mailbox}/Archive\";"
);
assert_eq!(renamed, ["vnd.stalwart.default_mailbox"]);
}
#[test]
fn renames_env_references_in_text_blocks() {
let src = "vacation text:\nHi ${env.vnd.stalwart.username}.\n.\n;\n";
let (out, _) = run(src).unwrap();
assert_eq!(
out,
"vacation text:\nHi ${env.vnd.inbuxa.username}.\n.\n;\n"
);
}
#[test]
fn leaves_unrelated_strings_comments_and_text_alone() {
let src = "# vnd.stalwart.while is old\n/* \"vnd.stalwart.x\" */\n\
if header :contains \"subject\" \"vnd.stalwart.while\" { fileinto \"vnd.stalwart.box\"; }\n\
vacation text:\nrequire \"vnd.stalwart.while\";\n.\n;\n";
assert!(run(src).is_none());
}
#[test]
fn require_context_ends_at_the_semicolon() {
let src = "require \"fileinto\"; fileinto \"vnd.stalwart.folder\";";
assert!(run(src).is_none());
}
#[test]
fn nothing_to_do_is_none() {
assert!(run("require \"fileinto\";\nkeep;\n").is_none());
assert!(rewrite(&[0xff, 0xfe, b'v']).is_none());
}
#[test]
fn target_detection() {
assert!(target_uses_inbuxa_names(&[
"fileinto".to_owned(),
"vnd.inbuxa.while".to_owned()
]));
assert!(!target_uses_inbuxa_names(&[
"fileinto".to_owned(),
"vnd.stalwart.while".to_owned()
]));
assert!(!target_uses_inbuxa_names(&[]));
}
}
+98 -7
View File
@@ -1,15 +1,16 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
use std::collections::HashSet;
use std::collections::{HashMap, HashSet};
use std::fmt::Write as _;
use serde_json::Value;
use serde_json::{Map, Value};
use super::common::{create_batch, jid, target_query_get};
use super::common::{create_batch, jid, target_query_get, update_batch};
use super::{Maps, Net, Plan, Uploader};
use crate::error::Error;
use crate::logging::Logger;
@@ -42,10 +43,10 @@ pub fn reconcile(
logger: &Logger,
) -> Result<Plan, Error> {
let targets = target_query_get(net, ty, None).map_err(Error::from)?;
let mut by_uid: std::collections::HashMap<String, String> = std::collections::HashMap::new();
let mut by_uid: HashMap<String, (String, &Value)> = HashMap::new();
for t in &targets {
if let (Some(uid), Some(id)) = (target_uid(t), jid(t)) {
by_uid.entry(uid).or_insert(id);
by_uid.entry(uid).or_insert((id, t));
}
}
@@ -76,12 +77,23 @@ pub fn reconcile(
};
let mut matched_uids: HashSet<String> = HashSet::new();
let mut updates: Vec<(String, Value)> = Vec::new();
let blobs = Uploader::new(net, &ctx.conn);
for (local, uid) in &rows {
if let Some(tid) = by_uid.get(uid) {
if let Some((tid, existing)) = by_uid.get(uid) {
maps.insert(ty, *local, crate::jmap::wire::JmapId(tid.clone()));
matched_uids.insert(uid.clone());
counts.skipped += 1;
match build_wire(ctx, ty, *local, maps, &blobs) {
Ok(wire) => match changed_properties(&wire, existing) {
Some(patch) => updates.push((tid.clone(), patch)),
None => counts.skipped += 1,
},
Err(e) if e.aborts_run() => return Err(e),
Err(e) => {
logger.warn(&format!("{} not compared: {e}", describe(ty, *local, uid)));
counts.skipped += 1;
}
}
continue;
}
let cid = format!("c{local}");
@@ -123,6 +135,8 @@ pub fn reconcile(
}
}
update_batch(net, ty, updates, counts, logger);
let objs: Vec<TargetObj> = targets
.iter()
.filter_map(|t| {
@@ -172,3 +186,80 @@ fn build_wire(
calendar_event_to_wire(&cal, dr != 0, ud != 0, &data, maps, blobs).map_err(Error::from)
}
}
/// An `updated` value as a point in time, so that offsets and fractional
/// seconds compare as the same instant. `None` if it does not parse.
fn parse_updated(s: &str) -> Option<time::OffsetDateTime> {
time::OffsetDateTime::parse(s, &time::format_description::well_known::Rfc3339).ok()
}
/// The update that makes `target` match the archive's `wire` object, or
/// `None` when nothing changed. When both carry `updated`, it decides: the
/// archive's copy wins only if it is newer. Otherwise each property the
/// archive writes is compared, and those that differ are sent whole.
fn changed_properties(wire: &Value, target: &Value) -> Option<Value> {
let wire = wire.as_object()?;
let stamp = |v: Option<&Value>| v.and_then(Value::as_str).and_then(parse_updated);
if let (Some(ours), Some(theirs)) = (stamp(wire.get("updated")), stamp(target.get("updated")))
&& ours <= theirs
{
return None;
}
let mut patch = Map::new();
for (k, v) in wire {
if k == "uid" || k == "id" {
continue;
}
if target.get(k) != Some(v) {
patch.insert(k.clone(), v.clone());
}
}
(!patch.is_empty()).then_some(Value::Object(patch))
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn unchanged_object_needs_no_update() {
let wire = json!({"uid": "u", "name": {"full": "Ann"}, "addressBookIds": {"A": true}});
let target = json!({"id": "T", "uid": "u", "name": {"full": "Ann"},
"addressBookIds": {"A": true}, "extra": 1});
assert_eq!(changed_properties(&wire, &target), None);
}
#[test]
fn changed_properties_are_sent_whole() {
let wire = json!({"uid": "u", "name": {"full": "Ann B"}, "addressBookIds": {"A": true}});
let target = json!({"id": "T", "uid": "u", "name": {"full": "Ann"},
"addressBookIds": {"A": true}});
assert_eq!(
changed_properties(&wire, &target),
Some(json!({"name": {"full": "Ann B"}}))
);
}
#[test]
fn updated_compares_instants_not_strings() {
let target = json!({"uid": "u", "title": "old", "updated": "2026-01-02T00:00:00Z"});
let same_instant = json!({"uid": "u", "title": "new",
"updated": "2026-01-02T01:00:00.000+01:00"});
assert_eq!(changed_properties(&same_instant, &target), None);
let later = json!({"uid": "u", "title": "new", "updated": "2026-01-02T00:00:00.5Z"});
assert!(changed_properties(&later, &target).is_some());
}
#[test]
fn updated_decides_when_both_sides_carry_it() {
let target = json!({"uid": "u", "title": "old", "updated": "2026-01-02T00:00:00Z"});
let older = json!({"uid": "u", "title": "new", "updated": "2026-01-01T00:00:00Z"});
assert_eq!(changed_properties(&older, &target), None, "target is newer");
let newer = json!({"uid": "u", "title": "new", "updated": "2026-01-03T00:00:00Z"});
assert_eq!(
changed_properties(&newer, &target),
Some(json!({"title": "new", "updated": "2026-01-03T00:00:00Z"}))
);
}
}
+3 -1
View File
@@ -1,5 +1,6 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -18,6 +19,7 @@ use crate::sync::{CommonConfig, RunOutcome, Summary, TypeCounts};
use super::collections;
use super::items;
use super::tree;
use crate::net::CertOverride;
#[derive(Debug, Clone, Copy)]
pub enum DavKindArg {
@@ -105,7 +107,7 @@ fn run_into(
let client = DavClient::new(
config.auth.to_jmap_auth(),
RetryPolicy::new(common.max_retries),
common.allow_invalid_certs,
CertOverride::for_url(common.allow_invalid_certs, &config.url),
);
client.set_logger(logger);
+58 -38
View File
@@ -20,6 +20,7 @@ use crate::sync::{CommonConfig, Summary, TypeCounts};
use super::folders::{self, plan_folders};
use super::{calendar, contacts, messages};
use crate::net::CertOverride;
#[derive(Debug, Clone)]
pub enum EwsAuth {
@@ -45,8 +46,8 @@ pub fn run(common: CommonConfig, config: EwsImportConfig) -> Result<Summary, Err
let logger = common.logger;
let mut conn = db::init::open(&common.archive)?;
let (auth, acquired) = resolve_auth(&config.auth, common.allow_invalid_certs)?;
let discovery = run_autodiscover(&config, &acquired, common.allow_invalid_certs)?;
let (auth, acquired) = resolve_auth(&config.auth)?;
let (discovery, certs) = run_autodiscover(&config, &acquired, common.allow_invalid_certs)?;
if logger.enabled(LEVEL_PROGRESS) {
eprintln!(
"EWS discovery: url={} source={:?}",
@@ -77,7 +78,7 @@ pub fn run(common: CommonConfig, config: EwsImportConfig) -> Result<Summary, Err
let client = EwsClient::new(
auth,
RetryPolicy::new(common.max_retries),
common.allow_invalid_certs,
certs.narrowed_to(&discovery.ews_url),
);
client.set_logger(logger);
if matches!(config.mailbox_kind, MailboxKind::PublicFolders) {
@@ -88,13 +89,7 @@ pub fn run(common: CommonConfig, config: EwsImportConfig) -> Result<Summary, Err
if let EwsAuth::OAuth(OAuthFlow::ClientCredentials { .. }) = &config.auth {
client.set_impersonation(Some(mailbox.clone()));
}
spawn_token_refresher(
&client,
&config.auth,
&acquired,
common.allow_invalid_certs,
logger,
);
spawn_token_refresher(&client, &config.auth, &acquired, logger);
let username = match &config.auth {
EwsAuth::Basic { user, .. } => user.clone(),
@@ -211,10 +206,7 @@ fn run_dry(
Ok(summary)
}
fn resolve_auth(
auth: &EwsAuth,
allow_invalid_certs: bool,
) -> Result<(Auth, Option<AcquiredToken>), Error> {
fn resolve_auth(auth: &EwsAuth) -> Result<(Auth, Option<AcquiredToken>), Error> {
match auth {
EwsAuth::Basic { user, password } => Ok((
Auth::Basic {
@@ -224,12 +216,9 @@ fn resolve_auth(
None,
)),
EwsAuth::Bearer { token } => {
let acq = acquire(
&OAuthFlow::PreAcquired {
token: token.clone(),
},
allow_invalid_certs,
)
let acq = acquire(&OAuthFlow::PreAcquired {
token: token.clone(),
})
.map_err(Error::from)?;
Ok((
Auth::Bearer {
@@ -239,7 +228,7 @@ fn resolve_auth(
))
}
EwsAuth::OAuth(flow) => {
let acq = acquire(flow, allow_invalid_certs).map_err(Error::from)?;
let acq = acquire(flow).map_err(Error::from)?;
Ok((
Auth::Bearer {
token: acq.access_token.clone(),
@@ -254,19 +243,36 @@ fn run_autodiscover(
config: &EwsImportConfig,
acquired: &Option<AcquiredToken>,
allow_invalid_certs: bool,
) -> Result<DiscoveryResult, Error> {
) -> Result<(DiscoveryResult, CertOverride), Error> {
let email = config
.mailbox
.clone()
.or_else(|| acquired.as_ref().and_then(|a| a.upn.clone()));
let result = discover(
config.url.as_deref(),
email.as_deref(),
None,
allow_invalid_certs,
)
.map_err(Error::from)?;
Ok(result)
let certs =
autodiscover_cert_override(config.url.as_deref(), email.as_deref(), allow_invalid_certs);
let result =
discover(config.url.as_deref(), email.as_deref(), None, &certs).map_err(Error::from)?;
Ok((result, certs))
}
/// Where `--allow-invalid-certs` applies for an EWS import: the host of
/// `--url` when one is given, and otherwise the mailbox's own domain, which is
/// where on-premises Autodiscover looks. Microsoft's hosts are never covered.
fn autodiscover_cert_override(
url: Option<&str>,
email: Option<&str>,
enabled: bool,
) -> CertOverride {
match (
url,
email
.and_then(|e| e.rsplit_once('@'))
.map(|(_, domain)| domain),
) {
(Some(url), _) => CertOverride::for_url(enabled, url),
(None, Some(domain)) => CertOverride::for_domain(enabled, domain),
(None, None) => CertOverride::none(),
}
}
fn resolve_mailbox(
@@ -370,7 +376,6 @@ fn spawn_token_refresher(
client: &EwsClient,
auth: &EwsAuth,
initial: &Option<AcquiredToken>,
allow_invalid_certs: bool,
logger: crate::logging::Logger,
) {
let flow = match auth {
@@ -402,14 +407,9 @@ fn spawn_token_refresher(
let result = if let (Some(rt), OAuthFlow::DeviceCode { tenant, client_id }) =
(refresh_token.as_deref(), &flow)
{
crate::exchange_ews::oauth::refresh_with_token(
tenant,
client_id,
rt,
allow_invalid_certs,
)
crate::exchange_ews::oauth::refresh_with_token(tenant, client_id, rt)
} else {
crate::exchange_ews::oauth::acquire(&flow, allow_invalid_certs)
crate::exchange_ews::oauth::acquire(&flow)
};
match result {
Ok(tok) => {
@@ -451,6 +451,26 @@ fn run_gc(conn: &Connection) -> Result<(), Error> {
mod tests {
use super::*;
#[test]
fn cert_override_follows_url_then_mailbox_domain() {
let by_url = autodiscover_cert_override(
Some("https://mail.corp.example/EWS/Exchange.asmx"),
Some("[email protected]"),
true,
);
assert!(by_url.allows("https://mail.corp.example/EWS/Exchange.asmx"));
assert!(!by_url.allows("https://autodiscover.corp.example/"));
let by_domain = autodiscover_cert_override(None, Some("[email protected]"), true);
assert!(
by_domain.allows("https://autodiscover.corp.example/autodiscover/autodiscover.xml")
);
assert!(!by_domain.allows("https://outlook.office365.com/EWS/Exchange.asmx"));
assert!(!autodiscover_cert_override(None, Some("[email protected]"), false).is_active());
assert!(!autodiscover_cert_override(None, None, true).is_active());
}
#[test]
fn synthetic_account_id_uses_smtp_for_primary() {
assert_eq!(
+7 -18
View File
@@ -20,6 +20,7 @@ use crate::exchange_graph::oauth::{
use crate::exchange_graph::types::{EventBodyFormat, MailboxKind, Surfaces, synthetic_account_id};
use crate::jmap::http::RetryPolicy;
use crate::logging::LEVEL_DEFAULT;
use crate::net::CertOverride;
use crate::sync::{CommonConfig, Summary, TypeCounts};
#[derive(Debug, Clone)]
@@ -68,11 +69,11 @@ pub fn run(common: CommonConfig, config: GraphImportConfig) -> Result<Summary, E
let logger = common.logger;
let mut conn = db::init::open(&common.archive)?;
let acquired = acquire_with_flow(&config.auth, common.allow_invalid_certs)?;
let acquired = acquire_with_flow(&config.auth)?;
let client = GraphClient::new(
acquired.access_token.clone(),
RetryPolicy::new(common.max_retries),
common.allow_invalid_certs,
CertOverride::for_url(common.allow_invalid_certs, &config.api_base),
);
client.set_logger(logger);
@@ -121,13 +122,7 @@ pub fn run(common: CommonConfig, config: GraphImportConfig) -> Result<Summary, E
&principal.user_principal_name,
)?;
let _refresher = spawn_token_refresher(
&client,
&config.auth,
&acquired,
common.allow_invalid_certs,
logger,
);
let _refresher = spawn_token_refresher(&client, &config.auth, &acquired, logger);
let mut summary = Summary::default();
let mut mailbox_counts = TypeCounts::default();
@@ -282,7 +277,7 @@ pub fn enumerate_mail_folders(
Ok(all)
}
fn acquire_with_flow(auth: &GraphAuth, allow_invalid_certs: bool) -> Result<AcquiredToken, Error> {
fn acquire_with_flow(auth: &GraphAuth) -> Result<AcquiredToken, Error> {
let flow = match auth {
GraphAuth::PreAcquired { token } => OAuthFlow::PreAcquired {
token: token.clone(),
@@ -295,7 +290,7 @@ fn acquire_with_flow(auth: &GraphAuth, allow_invalid_certs: bool) -> Result<Acqu
client_id: client_id.clone(),
},
};
acquire(&flow, allow_invalid_certs).map_err(Error::from)
acquire(&flow).map_err(Error::from)
}
fn resolve_endpoints(config: &GraphImportConfig, client: &GraphClient) -> Result<Endpoints, Error> {
@@ -379,7 +374,6 @@ fn spawn_token_refresher(
client: &GraphClient,
auth: &GraphAuth,
initial: &AcquiredToken,
allow_invalid_certs: bool,
logger: crate::logging::Logger,
) -> Option<TokenRefresher> {
let (authority, client_id) = match auth {
@@ -416,12 +410,7 @@ fn spawn_token_refresher(
break;
}
}
match refresh_access_token(
&authority,
&client_id,
&refresh_token,
allow_invalid_certs,
) {
match refresh_access_token(&authority, &client_id, &refresh_token) {
Ok(tok) => {
client.set_bearer(tok.access_token.clone());
if let Some(new_refresh) = tok.refresh_token {
+3 -1
View File
@@ -1,5 +1,6 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <hello@stalw.art>
* SPDX-FileCopyrightText: 2026 John Coffey <johnellis@linux.com>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -41,6 +42,7 @@ pub(crate) fn table_name(ty: ObjectType) -> &'static str {
use crate::jmap::account::AccountSelector;
use crate::jmap::http::{Auth, HttpClient, RetryPolicy};
use crate::logging::Logger;
use crate::net::CertOverride;
use crate::types::ObjectType;
pub struct CommonConfig {
@@ -134,7 +136,7 @@ impl Context {
let client = HttpClient::new(
connect.auth.clone(),
RetryPolicy::new(common.max_retries),
common.allow_invalid_certs,
CertOverride::for_url(common.allow_invalid_certs, &connect.url),
);
Ok(Context {
conn,
+2 -1
View File
@@ -11,6 +11,7 @@ mod seeder;
use inbuxa_migrate::jmap::account::{self, AccountSelector};
use inbuxa_migrate::jmap::http::{Auth, HttpClient, RetryPolicy};
use inbuxa_migrate::jmap::session::Session;
use inbuxa_migrate::net::CertOverride;
use integration::stalwart::shared as shared_stalwart;
fn admin_client() -> HttpClient {
@@ -20,7 +21,7 @@ fn admin_client() -> HttpClient {
password: seeder::ADMIN_PASSWORD.into(),
},
RetryPolicy::new(5),
true,
CertOverride::for_url(true, shared_stalwart().base_url()),
)
}
+2 -1
View File
@@ -13,6 +13,7 @@ use inbuxa_migrate::dav::parse::{parse_multistatus, strip_ascii_control_chars};
use inbuxa_migrate::dav::xml;
use inbuxa_migrate::jmap::error::JmapError;
use inbuxa_migrate::jmap::http::{Auth, RetryPolicy};
use inbuxa_migrate::net::CertOverride;
fn client(retries: u32) -> DavClient {
DavClient::new(
@@ -21,7 +22,7 @@ fn client(retries: u32) -> DavClient {
password: "p".into(),
},
RetryPolicy::new(retries),
false,
CertOverride::none(),
)
}
+3 -2
View File
@@ -22,6 +22,7 @@ use inbuxa_migrate::exchange_ews::xml::{
get_item_body, sync_folder_items_body,
};
use inbuxa_migrate::jmap::http::{Auth, RetryPolicy};
use inbuxa_migrate::net::CertOverride;
use mockito::Matcher;
const TXT_XML: &str = "text/xml; charset=utf-8";
@@ -33,7 +34,7 @@ fn client(retries: u32) -> EwsClient {
token: "t".to_owned(),
},
RetryPolicy::new(retries),
false,
CertOverride::none(),
)
}
@@ -63,7 +64,7 @@ fn autodiscover_v2_returns_global_endpoint() {
let _ = server;
let url = "https://outlook.office365.com/EWS/Exchange.asmx";
assert!(inbuxa_migrate::exchange_ews::autodiscover::is_fully_qualified_ews_url(url));
let r = discover(Some(url), None, None, false).unwrap();
let r = discover(Some(url), None, None, &CertOverride::none()).unwrap();
assert_eq!(r.source, DiscoverySource::SuppliedUrl);
assert_eq!(r.ews_url, url);
}
+11 -2
View File
@@ -21,6 +21,7 @@ use inbuxa_migrate::exchange_graph::recurrence::convert_patterned_recurrence;
use inbuxa_migrate::exchange_graph::retry::{HttpClass, classify_http_status};
use inbuxa_migrate::exchange_graph::types::Surfaces;
use inbuxa_migrate::jmap::http::RetryPolicy;
use inbuxa_migrate::net::CertOverride;
use mockito::{Matcher, Server};
use serde_json::json;
@@ -28,7 +29,11 @@ static INIT: Once = Once::new();
fn client_with_retries(retries: u32) -> GraphClient {
INIT.call_once(|| {});
GraphClient::new("BEARER".to_owned(), RetryPolicy::new(retries), false)
GraphClient::new(
"BEARER".to_owned(),
RetryPolicy::new(retries),
CertOverride::none(),
)
}
fn url_message_collection(server_url: &str, folder: &str, top: usize) -> String {
@@ -1681,7 +1686,11 @@ fn graph_client_retries_after_401_when_bearer_is_swapped() {
.expect(1)
.create();
let base = server.url();
let client = GraphClient::new("EXPIRED".to_owned(), RetryPolicy::new(0), false);
let client = GraphClient::new(
"EXPIRED".to_owned(),
RetryPolicy::new(0),
CertOverride::none(),
);
let url = format!("{base}/me");
let err = client.get(&url, Accept::Json).unwrap_err();
assert!(matches!(err, GraphError::Auth(_)));
+2 -1
View File
@@ -17,6 +17,7 @@ use inbuxa_migrate::jmap::session::{Limits, Session};
use inbuxa_migrate::jmap::wire::JmapId;
use inbuxa_migrate::jmap::wire::identity::Identity;
use inbuxa_migrate::jmap::wire::mailbox::Mailbox;
use inbuxa_migrate::net::CertOverride;
use serde_json::json;
fn client(retries: u32) -> HttpClient {
@@ -26,7 +27,7 @@ fn client(retries: u32) -> HttpClient {
password: "p".into(),
},
RetryPolicy::new(retries),
false,
CertOverride::none(),
)
}
+604 -5
View File
@@ -2575,7 +2575,7 @@ fn export_archive_read_failure_while_inlining_exits_seven() {
}
#[test]
fn export_sieve_script_matches_by_name_not_content() {
fn export_sieve_script_matched_by_name_is_updated_when_its_content_differs() {
let mut server = mockito::Server::new();
let base = server.url();
let api = "/jmap/api";
@@ -2615,9 +2615,125 @@ fn export_sieve_script_matches_by_name_not_content() {
)
.expect(1)
.create();
let no_download = server
let download = server
.mock("GET", Matcher::Regex("/jmap/dl/w/BSRV/.*".into()))
.with_body(b"unused".as_slice())
.with_body(b"keep;\n".as_slice())
.expect(1)
.create();
let upload = server
.mock("POST", Matcher::Regex("/jmap/upload/".into()))
.with_body(json!({"blobId":"UPN"}).to_string())
.expect(2)
.create();
let update = server
.mock("POST", api)
.match_body(Matcher::AllOf(vec![
Matcher::Regex("SieveScript/set".into()),
Matcher::Regex("\"update\":\\{\"S1\":\\{\"blobId\":\"UPN\"".into()),
]))
.with_body(
json!({"methodResponses":[["SieveScript/set",{"accountId":"w",
"updated":{"S1":null}},"s"]]})
.to_string(),
)
.expect(1)
.create();
let create = server
.mock("POST", api)
.match_body(Matcher::AllOf(vec![
Matcher::Regex("SieveScript/set".into()),
Matcher::Regex("reject".into()),
]))
.with_body(
json!({"methodResponses":[["SieveScript/set",{"accountId":"w",
"created":{"c2":{"id":"S2"}}},"s"]]})
.to_string(),
)
.expect(1)
.create();
let _activate = server
.mock("POST", api)
.match_body(Matcher::Regex("onSuccessActivateScript".into()))
.with_body(
json!({"methodResponses":[["SieveScript/set",{"accountId":"w"},"a"]]}).to_string(),
)
.expect(1)
.create();
let summary = sync::export::run(
common(&archive),
export_cfg_objects(&base, vec![ObjectType::SieveScript]),
)
.expect("export");
upload.assert();
create.assert();
download.assert();
update.assert();
let counts = summary
.per_type
.iter()
.find(|(t, _)| *t == "SieveScript")
.map(|(_, c)| c.clone())
.expect("sieve counts");
assert_eq!(
counts.updated, 1,
"the name-matched script gets the archive's content"
);
assert_eq!(counts.skipped, 0);
assert_eq!(counts.created, 1, "the unmatched name is created");
assert_eq!(counts.failed, 0);
let _ = std::fs::remove_file(&archive);
}
#[test]
fn export_sieve_script_matched_by_name_with_the_same_content_is_left_alone() {
let mut server = mockito::Server::new();
let base = server.url();
let api = "/jmap/api";
let archive = tmp();
let keepall_local = b"require [\"fileinto\"];\nkeep;\n";
let reject_local = b"require [\"reject\"];\nreject \"go away\";\n";
{
let conn = db::init::open(&archive).unwrap();
let blob1 = db::blobs::intern_blob(&conn, keepall_local).unwrap();
let blob2 = db::blobs::intern_blob(&conn, reject_local).unwrap();
conn.execute(
"INSERT INTO sieve_scripts (id,name,is_active,blob_id) VALUES (1,'keepall',1,?1)",
rusqlite::params![blob1],
)
.unwrap();
conn.execute(
"INSERT INTO sieve_scripts (id,name,is_active,blob_id) VALUES (2,'reject',0,?1)",
rusqlite::params![blob2],
)
.unwrap();
}
let _root = server.mock("GET", "/").with_status(404).create();
let _wk = server
.mock("GET", "/.well-known/jmap")
.with_body(session_body_full(&base))
.expect_at_least(1)
.create();
let _g = server
.mock("POST", api)
.match_body(Matcher::Regex("SieveScript/get".into()))
.with_body(
json!({"methodResponses":[["SieveScript/get",{"accountId":"w","list":[
{"id":"S1","name":"keepall","isActive":false,"blobId":"BSRV"}
],"notFound":[]},"g"]]})
.to_string(),
)
.expect(1)
.create();
let download = server
.mock("GET", Matcher::Regex("/jmap/dl/w/BSRV/.*".into()))
.with_body(keepall_local.as_slice())
.expect(1)
.create();
let no_update = server
.mock("POST", api)
.match_body(Matcher::Regex("\"update\"".into()))
.expect(0)
.create();
let upload = server
@@ -2654,7 +2770,8 @@ fn export_sieve_script_matches_by_name_not_content() {
.expect("export");
upload.assert();
create.assert();
no_download.assert();
download.assert();
no_update.assert();
let counts = summary
.per_type
.iter()
@@ -2663,13 +2780,213 @@ fn export_sieve_script_matches_by_name_not_content() {
.expect("sieve counts");
assert_eq!(
counts.skipped, 1,
"name-matched script is skipped even though its content differs from the target"
"same name and same content: nothing to do"
);
assert_eq!(counts.updated, 0);
assert_eq!(counts.created, 1, "the unmatched name is created");
assert_eq!(counts.failed, 0);
let _ = std::fs::remove_file(&archive);
}
/// A session like `session_body_full`, whose Sieve capability lists
/// `extensions` in `sieveExtensions`.
fn session_body_sieve(base: &str, extensions: &[&str]) -> String {
let mut v: serde_json::Value = serde_json::from_str(&session_body_full(base)).unwrap();
v["accounts"]["w"]["accountCapabilities"]["urn:ietf:params:jmap:sieve"] =
json!({ "sieveExtensions": extensions });
v.to_string()
}
/// One active script `name` in a fresh archive, with `body` as its content.
fn archive_with_active_sieve(name: &str, body: &[u8]) -> PathBuf {
let archive = tmp();
let conn = db::init::open(&archive).unwrap();
let blob = db::blobs::intern_blob(&conn, body).unwrap();
conn.execute(
"INSERT INTO sieve_scripts (id,name,is_active,blob_id) VALUES (1,?1,1,?2)",
rusqlite::params![name, blob],
)
.unwrap();
archive
}
/// Mocks for exporting one script to an empty target: get, upload (matched
/// by `upload_body`), create, and activation answered with `activation`.
/// Returns the upload mock and the activation mock.
fn mock_sieve_export(
server: &mut mockito::ServerGuard,
session: String,
upload_body: Matcher,
create_ok: bool,
activation: serde_json::Value,
) -> (mockito::Mock, mockito::Mock, Vec<mockito::Mock>) {
let api = "/jmap/api";
let mut keep = vec![
server.mock("GET", "/").with_status(404).create(),
server
.mock("GET", "/.well-known/jmap")
.with_body(session)
.expect_at_least(1)
.create(),
server
.mock("POST", api)
.match_body(Matcher::Regex("SieveScript/get".into()))
.with_body(
json!({"methodResponses":[["SieveScript/get",
{"accountId":"w","list":[],"notFound":[]},"g"]]})
.to_string(),
)
.create(),
];
let upload = server
.mock("POST", Matcher::Regex("/jmap/upload/".into()))
.match_body(upload_body)
.with_body(json!({"blobId":"UPN"}).to_string())
.expect(1)
.create();
let created = if create_ok {
json!({"accountId":"w","created":{"c1":{"id":"S1"}}})
} else {
json!({"accountId":"w","notCreated":{"c1":{"type":"invalidScript",
"description":"unknown extension"}}})
};
keep.push(
server
.mock("POST", api)
.match_body(Matcher::AllOf(vec![
Matcher::Regex("SieveScript/set".into()),
Matcher::Regex("\"create\"".into()),
]))
.with_body(json!({"methodResponses":[["SieveScript/set", created, "s"]]}).to_string())
.create(),
);
let activate = server
.mock("POST", api)
.match_body(Matcher::AllOf(vec![
Matcher::Regex("SieveScript/set".into()),
Matcher::Regex("onSuccess".into()),
]))
.with_body(json!({"methodResponses":[activation]}).to_string())
.expect(1)
.create();
(upload, activate, keep)
}
fn sieve_counts(summary: &sync::Summary) -> sync::TypeCounts {
summary
.per_type
.iter()
.find(|(t, _)| *t == "SieveScript")
.map(|(_, c)| c.clone())
.expect("sieve counts")
}
#[test]
fn export_sieve_renames_stalwart_names_for_an_inbuxa_target() {
let mut server = mockito::Server::new();
let base = server.url();
let archive = archive_with_active_sieve(
"loop",
b"require [\"fileinto\", \"vnd.stalwart.while\"];\nkeep;\n",
);
let (upload, activate, _keep) = mock_sieve_export(
&mut server,
session_body_sieve(
&base,
&["fileinto", "vnd.inbuxa.while", "vnd.inbuxa.expressions"],
),
Matcher::Exact("require [\"fileinto\", \"vnd.inbuxa.while\"];\nkeep;\n".into()),
true,
json!(["SieveScript/set", {"accountId":"w"}, "a"]),
);
let summary = sync::export::run(
common(&archive),
export_cfg_objects(&base, vec![ObjectType::SieveScript]),
)
.expect("export");
upload.assert();
activate.assert();
let c = sieve_counts(&summary);
assert_eq!((c.created, c.failed), (1, 0));
let _ = std::fs::remove_file(&archive);
}
#[test]
fn export_sieve_keeps_stalwart_names_for_a_target_without_inbuxa_names() {
let mut server = mockito::Server::new();
let base = server.url();
let archive = archive_with_active_sieve(
"loop",
b"require [\"fileinto\", \"vnd.stalwart.while\"];\nkeep;\n",
);
let (upload, _activate, _keep) = mock_sieve_export(
&mut server,
session_body_sieve(&base, &["fileinto", "vnd.stalwart.while"]),
Matcher::Regex("vnd\\.stalwart\\.while".into()),
true,
json!(["SieveScript/set", {"accountId":"w"}, "a"]),
);
let summary = sync::export::run(
common(&archive),
export_cfg_objects(&base, vec![ObjectType::SieveScript]),
)
.expect("export");
upload.assert();
assert_eq!(sieve_counts(&summary).failed, 0);
let _ = std::fs::remove_file(&archive);
}
#[test]
fn export_sieve_activation_error_is_a_failure() {
let mut server = mockito::Server::new();
let base = server.url();
let archive = archive_with_active_sieve("main", b"require [\"fileinto\"];\nkeep;\n");
let (_upload, activate, _keep) = mock_sieve_export(
&mut server,
session_body_full(&base),
Matcher::Any,
true,
json!(["error", {"type":"invalidArguments","description":"cannot activate"}, "a"]),
);
let summary = sync::export::run(
common(&archive),
export_cfg_objects(&base, vec![ObjectType::SieveScript]),
)
.expect("export");
activate.assert();
let c = sieve_counts(&summary);
assert_eq!(c.created, 1);
assert_eq!(
c.failed, 1,
"a failed activation is a failure, not a warning"
);
assert!(summary.any_failed(), "so export exits non-zero");
let _ = std::fs::remove_file(&archive);
}
#[test]
fn export_sieve_active_script_not_created_is_a_failure() {
let mut server = mockito::Server::new();
let base = server.url();
let archive = archive_with_active_sieve("main", b"require [\"nope\"];\nkeep;\n");
let (_upload, _activate, _keep) = mock_sieve_export(
&mut server,
session_body_full(&base),
Matcher::Any,
false,
json!(["SieveScript/set", {"accountId":"w"}, "a"]),
);
let summary = sync::export::run(
common(&archive),
export_cfg_objects(&base, vec![ObjectType::SieveScript]),
)
.expect("export");
let c = sieve_counts(&summary);
assert_eq!((c.created, c.failed), (0, 1));
assert!(summary.any_failed());
let _ = std::fs::remove_file(&archive);
}
#[test]
fn export_sieve_scripts_identical_content_different_names_both_created() {
let mut server = mockito::Server::new();
@@ -4701,3 +5018,285 @@ fn export_different_messages_sharing_a_message_id_are_not_merged() {
no_set.assert();
let _ = std::fs::remove_file(&archive);
}
#[test]
fn export_rerun_carries_a_read_flag_set_at_the_source() {
let mut server = mockito::Server::new();
let base = server.url();
let api = "/jmap/api";
let archive = tmp();
let _setup = two_folder_archive_and_target(&mut server, &archive);
insert_email_copy(&archive, ONE_MESSAGE, 1);
{
let conn = db::init::open(&archive).unwrap();
conn.execute("UPDATE emails SET keywords='[\"$seen\"]'", [])
.unwrap();
}
let _eq = server
.mock("POST", api)
.match_body(Matcher::Regex("Email/query".into()))
.with_body(
json!({"methodResponses":[["Email/query",{"accountId":"w","ids":["X1"]},"q"]]})
.to_string(),
)
.expect(1)
.create();
let _eg = server
.mock("POST", api)
.match_body(Matcher::Regex("Email/get".into()))
.with_body(
json!({"methodResponses":[["Email/get",{"accountId":"w","list":[
{"id":"X1","messageId":["both@h"],"size":ONE_MESSAGE.len(),
"mailboxIds":{"T1":true},"keywords":{"$flagged":true}}
],"notFound":[]},"g"]]})
.to_string(),
)
.expect(1)
.create();
let set = server
.mock("POST", api)
.match_body(Matcher::AllOf(vec![
Matcher::Regex("Email/set".into()),
Matcher::Regex("\"keywords/\\$seen\":true".into()),
Matcher::Regex("\"keywords/\\$flagged\":null".into()),
]))
.with_body(
json!({"methodResponses":[["Email/set",{"accountId":"w","updated":{"X1":null}},"s"]]})
.to_string(),
)
.expect(1)
.create();
let summary = sync::export::run(
common(&archive),
export_cfg_objects(&base, vec![ObjectType::Mailbox, ObjectType::Email]),
)
.expect("export");
let email = email_counts(&summary);
assert_eq!(email.updated, 1);
assert_eq!(email.created, 0);
assert_eq!(email.failed, 0);
set.assert();
let _ = std::fs::remove_file(&archive);
}
fn contact_rerun(local_updated: &str, updates_sent: usize) -> inbuxa_migrate::sync::TypeCounts {
let mut server = mockito::Server::new();
let base = server.url();
let api = "/jmap/api";
let archive = tmp();
{
let conn = db::init::open(&archive).unwrap();
conn.execute(
"INSERT INTO address_books (id,name,description,is_default) VALUES (1,'Personal',NULL,1)",
[],
)
.unwrap();
let card = json!({"@type":"Card","version":"1.0","uid":"u1",
"name":{"full":"Ann Brown"},"updated":local_updated});
conn.execute(
"INSERT INTO contact_cards (id,uid,address_book_ids,data) VALUES (1,'u1','[1]',?1)",
rusqlite::params![card.to_string()],
)
.unwrap();
}
let _root = server.mock("GET", "/").with_status(404).create();
let _wk = server
.mock("GET", "/.well-known/jmap")
.with_body(session_body_full(&base))
.expect_at_least(1)
.create();
let _abg = server
.mock("POST", api)
.match_body(Matcher::Regex("AddressBook/get".into()))
.with_body(
json!({"methodResponses":[["AddressBook/get",{"accountId":"w","list":[
{"id":"P","name":"Personal","isDefault":true,"myRights":{"mayDelete":false}}
],"notFound":[]},"g"]]})
.to_string(),
)
.expect_at_least(1)
.create();
let _term = anchor_terminator(&mut server, api, "ContactCard");
let _cq = server
.mock("POST", api)
.match_body(Matcher::Regex("ContactCard/query".into()))
.with_body(
json!({"methodResponses":[["ContactCard/query",{"accountId":"w","ids":["C1"]},"q"]]})
.to_string(),
)
.expect(1)
.create();
let _cg = server
.mock("POST", api)
.match_body(Matcher::Regex("ContactCard/get".into()))
.with_body(
json!({"methodResponses":[["ContactCard/get",{"accountId":"w","list":[
{"id":"C1","@type":"Card","version":"1.0","uid":"u1",
"name":{"full":"Ann"},"addressBookIds":{"P":true},
"updated":"2026-02-01T00:00:00Z"}
],"notFound":[]},"g"]]})
.to_string(),
)
.expect(1)
.create();
let update = server
.mock("POST", api)
.match_body(Matcher::AllOf(vec![
Matcher::Regex("ContactCard/set".into()),
Matcher::Regex("\"update\"".into()),
Matcher::Regex("Ann Brown".into()),
]))
.with_body(
json!({"methodResponses":[["ContactCard/set",{"accountId":"w","updated":{"C1":null}},"s"]]})
.to_string(),
)
.expect(updates_sent)
.create();
let summary = sync::export::run(
common(&archive),
export_cfg_objects(
&base,
vec![ObjectType::AddressBook, ObjectType::ContactCard],
),
)
.expect("export");
let counts = summary
.per_type
.iter()
.find(|(t, _)| *t == "ContactCard")
.map(|(_, c)| c.clone())
.expect("contact counts");
update.assert();
let _ = std::fs::remove_file(&archive);
counts
}
#[test]
fn export_rerun_updates_a_contact_edited_at_the_source() {
let counts = contact_rerun("2026-03-01T00:00:00Z", 1);
assert_eq!(counts.updated, 1, "the newer archive copy is written");
assert_eq!(counts.failed, 0);
}
#[test]
fn export_rerun_leaves_a_contact_alone_when_the_target_is_newer() {
let counts = contact_rerun("2026-01-01T00:00:00Z", 0);
assert_eq!(
counts.updated, 0,
"the target's copy is newer, so nothing is sent"
);
assert_eq!(counts.skipped, 1);
}
fn event_rerun(local_updated: &str, updates_sent: usize) -> inbuxa_migrate::sync::TypeCounts {
let mut server = mockito::Server::new();
let base = server.url();
let api = "/jmap/api";
let archive = tmp();
{
let conn = db::init::open(&archive).unwrap();
conn.execute(
"INSERT INTO calendars (id,name,is_default) VALUES (1,'Work',1)",
[],
)
.unwrap();
let event = json!({"@type":"Event","uid":"ev1","title":"Planning, moved",
"start":"2026-03-02T10:00:00","duration":"PT1H",
"updated":local_updated});
conn.execute(
"INSERT INTO calendar_events (id,calendar_ids,is_draft,use_default_alerts,data)
VALUES (1,'[1]',0,0,?1)",
rusqlite::params![event.to_string()],
)
.unwrap();
}
let _root = server.mock("GET", "/").with_status(404).create();
let _wk = server
.mock("GET", "/.well-known/jmap")
.with_body(session_body_full(&base))
.expect_at_least(1)
.create();
let _calg = server
.mock("POST", api)
.match_body(Matcher::Regex("Calendar/get".into()))
.with_body(
json!({"methodResponses":[["Calendar/get",{"accountId":"w","list":[
{"id":"K","name":"Work","isDefault":true,"myRights":{"mayDelete":false}}
],"notFound":[]},"g"]]})
.to_string(),
)
.expect_at_least(1)
.create();
let _term = anchor_terminator(&mut server, api, "CalendarEvent");
let _eq = server
.mock("POST", api)
.match_body(Matcher::Regex("CalendarEvent/query".into()))
.with_body(
json!({"methodResponses":[["CalendarEvent/query",{"accountId":"w","ids":["E1"]},"q"]]})
.to_string(),
)
.expect(1)
.create();
let _eg = server
.mock("POST", api)
.match_body(Matcher::Regex("CalendarEvent/get".into()))
.with_body(
json!({"methodResponses":[["CalendarEvent/get",{"accountId":"w","list":[
{"id":"E1","@type":"Event","uid":"ev1","title":"Planning",
"start":"2026-03-01T10:00:00","duration":"PT1H",
"calendarIds":{"K":true},"isDraft":false,"useDefaultAlerts":false,
"updated":"2026-02-01T00:00:00Z"}
],"notFound":[]},"g"]]})
.to_string(),
)
.expect(1)
.create();
let update = server
.mock("POST", api)
.match_body(Matcher::AllOf(vec![
Matcher::Regex("CalendarEvent/set".into()),
Matcher::Regex("\"update\"".into()),
Matcher::Regex("Planning, moved".into()),
]))
.with_body(
json!({"methodResponses":[["CalendarEvent/set",{"accountId":"w","updated":{"E1":null}},"s"]]})
.to_string(),
)
.expect(updates_sent)
.create();
let summary = sync::export::run(
common(&archive),
export_cfg_objects(&base, vec![ObjectType::Calendar, ObjectType::CalendarEvent]),
)
.expect("export");
let counts = summary
.per_type
.iter()
.find(|(t, _)| *t == "CalendarEvent")
.map(|(_, c)| c.clone())
.expect("event counts");
update.assert();
let _ = std::fs::remove_file(&archive);
counts
}
#[test]
fn export_rerun_updates_an_event_moved_at_the_source() {
let counts = event_rerun("2026-03-01T00:00:00Z", 1);
assert_eq!(counts.updated, 1, "the newer archive copy is written");
assert_eq!(counts.failed, 0);
}
#[test]
fn export_rerun_leaves_an_event_alone_when_nothing_is_newer() {
let counts = event_rerun("2026-02-01T01:00:00+01:00", 0);
assert_eq!(
counts.updated, 0,
"same instant as the target's, written another way"
);
assert_eq!(counts.skipped, 1);
}
+22 -5
View File
@@ -16,6 +16,7 @@ use inbuxa_migrate::jmap::http::{Auth, HttpClient, RetryPolicy};
use inbuxa_migrate::jmap::request::Request;
use inbuxa_migrate::jmap::session::Session;
use inbuxa_migrate::logging::Logger;
use inbuxa_migrate::net::CertOverride;
use inbuxa_migrate::sync::{self, CommonConfig, ConnectConfig, ExportConfig, ImportConfig};
use integration::stalwart::shared as shared_stalwart;
use rusqlite::Connection;
@@ -785,7 +786,11 @@ fn export_inlines_contact_and_event_blobs_instead_of_blob_ids() {
assert_eq!(counts.failed, 0, "{name} had no failures: {counts:?}");
}
let client = HttpClient::new(basic("test6"), RetryPolicy::new(5), true);
let client = HttpClient::new(
basic("test6"),
RetryPolicy::new(5),
CertOverride::for_url(true, base_url()),
);
let session = Session::discover(&client, base_url()).expect("discover target session");
let api = session.api_url.clone();
@@ -1047,7 +1052,11 @@ fn live_burst_exceeds_concurrent_requests_and_recovers() {
let fx = seeder::provision(base_url()).expect("provision");
let acc = fx.account("test1").expect("test1");
let client = HttpClient::new(basic("test1"), RetryPolicy::new(20), true);
let client = HttpClient::new(
basic("test1"),
RetryPolicy::new(20),
CertOverride::for_url(true, base_url()),
);
let session = Session::discover(&client, base_url()).expect("discover session");
let server_limits = session.core_limits().expect("core limits");
@@ -1148,7 +1157,7 @@ impl JmapSettingsGuard {
password: seeder::ADMIN_PASSWORD.to_owned(),
},
RetryPolicy::new(5),
true,
CertOverride::for_url(true, base_url()),
);
let session = Session::discover(&admin, base_url()).expect("admin discover");
let admin_account = session
@@ -1269,7 +1278,11 @@ fn live_blob_quota_429_triggers_retry_after_then_succeeds() {
);
let _ttl_guard = JmapSettingsGuard::override_settings(updates);
let client = HttpClient::new(basic("test1"), RetryPolicy::new(20), true);
let client = HttpClient::new(
basic("test1"),
RetryPolicy::new(20),
CertOverride::for_url(true, base_url()),
);
let session = Session::discover(&client, base_url()).expect("discover session");
let limits = session.core_limits().expect("core limits");
client.set_limits(&limits);
@@ -1346,7 +1359,11 @@ fn import_delta_propagates_email_keyword_change_via_changes() {
.expect("an unflagged email exists in the archive")
};
let client = HttpClient::new(basic("test1"), RetryPolicy::new(5), true);
let client = HttpClient::new(
basic("test1"),
RetryPolicy::new(5),
CertOverride::for_url(true, base_url()),
);
let session = Session::discover(&client, base_url()).expect("session discovered");
let account = account::resolve(
&AccountSelector::Id(acc.account_id.clone()),