Compare commits

..
1 Commits
Author SHA1 Message Date
jcoffey-dev db76a1044f 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 2m32s
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; 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,
  and replaced with the archive's when it differs.

Updated items are counted as `updated`; unchanged ones stay `skipped`. The
usage guide now describes this.
2026-09-30 11:32:04 -07:00
29 changed files with 265 additions and 1641 deletions
+1 -20
View File
@@ -26,26 +26,7 @@ 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 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.
| `--allow-invalid-certs` | Accept self-signed / invalid TLS certs. |
Secrets come from the `INBUXA_MIGRATE_*` environment variables or a prompt;
see [Credentials](../README.md#credentials). The command line takes them too,
+1 -4
View File
@@ -122,10 +122,7 @@ struct GlobalArgs {
)]
max_retries: u32,
#[arg(
long,
help = "Accept an invalid TLS certificate from the --url host only; sign-in endpoints are always verified"
)]
#[arg(long, help = "Accept self-signed / invalid TLS certificates")]
allow_invalid_certs: bool,
}
+19 -51
View File
@@ -14,6 +14,7 @@ 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};
@@ -21,7 +22,6 @@ 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,8 +48,6 @@ pub struct MultiStatus {
struct Inner {
agent: Agent,
lax_agent: Option<Agent>,
certs: CertOverride,
auth: Auth,
retry: RetryPolicy,
rate_limit: RateLimitState,
@@ -59,42 +57,31 @@ 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, 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))
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(),
)
.build();
config.new_agent()
};
let lax_agent = certs.is_active().then(|| build(true));
DavClient {
inner: Arc::new(Inner {
agent: build(false),
lax_agent,
certs,
agent: config.new_agent(),
auth,
retry,
rate_limit: RateLimitState::new(),
@@ -835,7 +822,7 @@ impl DavClient {
let request = builder
.body(payload)
.map_err(|e| ureq::Error::Other(Box::new(std::io::Error::other(e))))?;
self.inner.agent_for(req.url).run(request)
self.inner.agent.run(request)
}
}
@@ -954,25 +941,6 @@ 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() {
@@ -982,7 +950,7 @@ mod tests {
password: "p".into(),
},
RetryPolicy::new(3),
CertOverride::none(),
false,
);
assert_eq!(c.retries_observed(), 0);
assert_eq!(c.retry_after_sleeps(), 0);
@@ -995,7 +963,7 @@ mod tests {
token: "abc".into(),
},
RetryPolicy::new(0),
CertOverride::none(),
false,
);
let logger = c.logger();
assert_eq!(logger.level(), LEVEL_DEFAULT);
-95
View File
@@ -1,6 +1,5 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -10,7 +9,6 @@ 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)?;
@@ -74,50 +72,6 @@ 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")?;
@@ -129,53 +83,4 @@ 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());
}
}
+19 -23
View File
@@ -1,6 +1,5 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -10,10 +9,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 =
@@ -49,7 +48,7 @@ pub fn discover(
supplied_url: Option<&str>,
email: Option<&str>,
auth_header: Option<&str>,
certs: &CertOverride,
allow_invalid_certs: bool,
) -> Result<DiscoveryResult, EwsError> {
if let Some(url) = supplied_url
&& is_fully_qualified_ews_url(url)
@@ -64,17 +63,8 @@ pub fn discover(
"either a fully-qualified --url or --mailbox is required".to_owned(),
));
};
// 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) {
let agent = build_agent(allow_invalid_certs);
if let Ok(url) = autodiscover_v2(&agent, email) {
return Ok(DiscoveryResult {
ews_url: url,
source: DiscoverySource::V2,
@@ -92,7 +82,7 @@ pub fn discover(
let candidates = pox_candidates(domain);
for candidate in &candidates {
tried.push(candidate.clone());
match autodiscover_v1(pick(candidate), candidate, &current_email, auth_header) {
match autodiscover_v1(&agent, candidate, &current_email, auth_header) {
Ok(PoxOutcome::EwsUrl(url)) => {
return Ok(DiscoveryResult {
ews_url: url,
@@ -114,7 +104,7 @@ pub fn discover(
url_redirects += 1;
tried.push(url.clone());
if let Ok(PoxOutcome::EwsUrl(u)) =
autodiscover_v1(pick(&url), &url, &current_email, auth_header)
autodiscover_v1(&agent, &url, &current_email, auth_header)
{
return Ok(DiscoveryResult {
ews_url: u,
@@ -140,13 +130,19 @@ pub fn discover(
)))
}
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();
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();
config.new_agent()
}
+18 -47
View File
@@ -12,6 +12,7 @@ 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};
@@ -21,15 +22,12 @@ 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>>,
@@ -44,17 +42,6 @@ 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>,
@@ -67,23 +54,23 @@ pub struct SoapResponse {
}
impl EwsClient {
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))
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(),
)
.build();
config.new_agent()
};
let lax_agent = certs.is_active().then(|| build(true));
EwsClient {
inner: Arc::new(Inner {
agent: build(false),
lax_agent,
certs,
agent: config.new_agent(),
auth: Mutex::new(auth),
impersonated_smtp: Mutex::new(None),
anchor_mailbox: Mutex::new(None),
@@ -421,7 +408,7 @@ impl EwsClient {
fn one_attempt(&self, url: &str, body: &str, action: &str) -> AttemptOutcome {
let mut req = self
.inner
.agent_for(url)
.agent
.post(url)
.header("Authorization", self.auth_header())
.header("Content-Type", "text/xml; charset=utf-8")
@@ -548,29 +535,13 @@ 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),
CertOverride::none(),
false,
);
assert_eq!(c.server_version(), ServerVersion::Exchange2013Sp1);
assert_eq!(c.retries_observed(), 0);
@@ -582,7 +553,7 @@ mod tests {
let c = EwsClient::new(
Auth::Bearer { token: "t".into() },
RetryPolicy::new(0),
CertOverride::none(),
false,
);
c.set_server_version(ServerVersion::Exchange2019);
assert_eq!(c.server_version(), ServerVersion::Exchange2019);
@@ -593,7 +564,7 @@ mod tests {
let c = EwsClient::new(
Auth::Bearer { token: "t".into() },
RetryPolicy::new(0),
CertOverride::none(),
false,
);
c.set_anchor_mailbox(Some("alice@x".to_owned()));
assert_eq!(c.anchor_header().as_deref(), Some("alice@x"));
+35 -19
View File
@@ -1,6 +1,5 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -11,9 +10,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 =
@@ -47,7 +46,7 @@ pub enum OAuthFlow {
},
}
pub fn acquire(flow: &OAuthFlow) -> Result<AcquiredToken, EwsError> {
pub fn acquire(flow: &OAuthFlow, allow_invalid_certs: bool) -> Result<AcquiredToken, EwsError> {
match flow {
OAuthFlow::PreAcquired { token } => {
let claims = decode_jwt_claims(token).unwrap_or_default();
@@ -64,8 +63,10 @@ pub fn acquire(flow: &OAuthFlow) -> Result<AcquiredToken, EwsError> {
tenant,
client_id,
client_secret,
} => client_credentials(tenant, client_id, client_secret),
OAuthFlow::DeviceCode { tenant, client_id } => device_code_flow(tenant, client_id),
} => client_credentials(tenant, client_id, client_secret, allow_invalid_certs),
OAuthFlow::DeviceCode { tenant, client_id } => {
device_code_flow(tenant, client_id, allow_invalid_certs)
}
}
}
@@ -103,13 +104,19 @@ fn device_code_endpoint(tenant: &str) -> String {
format!("https://login.microsoftonline.com/{tenant}/oauth2/v2.0/devicecode")
}
fn build_agent() -> ureq::Agent {
let config: Config = with_timeouts!(
Config::builder()
.http_status_as_error(false)
.tls_config(tls(false))
)
.build();
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();
config.new_agent()
}
@@ -117,8 +124,9 @@ fn client_credentials(
tenant: &str,
client_id: &str,
client_secret: &str,
allow_invalid_certs: bool,
) -> Result<AcquiredToken, EwsError> {
let agent = build_agent();
let agent = build_agent(allow_invalid_certs);
let body = form_encode(&[
("client_id", client_id),
("client_secret", client_secret),
@@ -134,8 +142,12 @@ fn client_credentials(
parse_token_response(resp)
}
fn device_code_flow(tenant: &str, client_id: &str) -> Result<AcquiredToken, EwsError> {
let agent = build_agent();
fn device_code_flow(
tenant: &str,
client_id: &str,
allow_invalid_certs: bool,
) -> Result<AcquiredToken, EwsError> {
let agent = build_agent(allow_invalid_certs);
let body = form_encode(&[("client_id", client_id), ("scope", SCOPE_DELEGATED)]);
let endpoint = device_code_endpoint(tenant);
let mut resp = agent
@@ -267,8 +279,9 @@ pub fn refresh_with_token(
tenant: &str,
client_id: &str,
refresh_token: &str,
allow_invalid_certs: bool,
) -> Result<AcquiredToken, EwsError> {
let agent = build_agent();
let agent = build_agent(allow_invalid_certs);
let body = form_encode(&[
("client_id", client_id),
("grant_type", "refresh_token"),
@@ -353,9 +366,12 @@ 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(),
})
let acq = acquire(
&OAuthFlow::PreAcquired {
token: token.clone(),
},
false,
)
.unwrap();
assert_eq!(acq.access_token, token);
assert_eq!(acq.tenant_id.as_deref(), Some("t-2"));
+17 -50
View File
@@ -13,6 +13,7 @@ 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;
@@ -20,7 +21,6 @@ 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,8 +63,6 @@ impl GraphResponse {
struct Inner {
agent: Agent,
lax_agent: Option<Agent>,
certs: CertOverride,
bearer: Mutex<String>,
retry: RetryPolicy,
rate_limit: RateLimitState,
@@ -75,17 +73,6 @@ 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>,
@@ -102,23 +89,23 @@ enum Attempt {
}
impl GraphClient {
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))
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(),
)
.build();
config.new_agent()
};
let lax_agent = certs.is_active().then(|| build(true));
GraphClient {
inner: Arc::new(Inner {
agent: build(false),
lax_agent,
certs,
agent: config.new_agent(),
bearer: Mutex::new(bearer),
retry,
rate_limit: RateLimitState::new(),
@@ -323,7 +310,7 @@ impl GraphClient {
extra_prefer: &[&str],
) -> Attempt {
let mut req = match method {
"GET" => self.inner.agent_for(url).get(url),
"GET" => self.inner.agent.get(url),
other => {
return Attempt::Transport(GraphError::Connect(format!(
"unsupported method {other} (graph importer is read-only)"
@@ -487,30 +474,10 @@ 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),
CertOverride::none(),
);
let c = GraphClient::new("token".to_owned(), RetryPolicy::new(3), false);
assert_eq!(c.retries_observed(), 0);
assert_eq!(c.retry_after_sleeps(), 0);
assert_eq!(c.requests_observed(), 0);
@@ -519,7 +486,7 @@ mod tests {
#[test]
fn bearer_can_be_swapped_at_runtime() {
let c = GraphClient::new("old".to_owned(), RetryPolicy::new(0), CertOverride::none());
let c = GraphClient::new("old".to_owned(), RetryPolicy::new(0), false);
c.set_bearer("new".to_owned());
assert_eq!(c.auth_header(), "Bearer new");
}
+24 -14
View File
@@ -1,6 +1,5 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -11,9 +10,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";
@@ -79,23 +78,29 @@ pub struct AcquiredToken {
pub name: Option<String>,
}
fn build_agent() -> ureq::Agent {
let config: Config = with_timeouts!(
Config::builder()
.http_status_as_error(false)
.tls_config(tls(false))
)
.build();
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();
config.new_agent()
}
pub fn acquire(flow: &OAuthFlow) -> Result<AcquiredToken, GraphError> {
pub fn acquire(flow: &OAuthFlow, allow_invalid_certs: bool) -> 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),
} => device_code_flow(authority, client_id, allow_invalid_certs),
}
}
@@ -208,8 +213,12 @@ pub fn parse_token_response(status: u16, json: &Value) -> TokenResponse {
}
}
fn device_code_flow(authority: &str, client_id: &str) -> Result<AcquiredToken, GraphError> {
let agent = build_agent();
fn device_code_flow(
authority: &str,
client_id: &str,
allow_invalid_certs: bool,
) -> Result<AcquiredToken, GraphError> {
let agent = build_agent(allow_invalid_certs);
let body = form_encode(&[("client_id", client_id), ("scope", SCOPES)]);
let endpoint = device_code_endpoint(authority);
let mut resp = agent
@@ -285,8 +294,9 @@ pub fn refresh_access_token(
authority: &str,
client_id: &str,
refresh_token: &str,
allow_invalid_certs: bool,
) -> Result<AcquiredToken, GraphError> {
let agent = build_agent();
let agent = build_agent(allow_invalid_certs);
let body = form_encode(&[
("client_id", client_id),
("grant_type", "refresh_token"),
+1 -3
View File
@@ -1,6 +1,5 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -171,7 +170,6 @@ 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!({
@@ -188,7 +186,7 @@ mod tests {
HttpClient::new(
crate::jmap::http::Auth::Bearer { token: "t".into() },
crate::jmap::http::RetryPolicy::new(0),
CertOverride::none(),
false,
)
}
+29 -102
View File
@@ -1,6 +1,5 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -14,6 +13,7 @@ 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,7 +21,6 @@ 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;
@@ -70,10 +69,9 @@ 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>,
@@ -83,17 +81,6 @@ 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,
@@ -116,25 +103,26 @@ enum Attempt {
}
impl HttpClient {
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))
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(),
)
.build();
config.new_agent()
};
let lax_agent = certs.is_active().then(|| build(true));
HttpClient {
inner: Arc::new(Inner {
agent: build(false),
lax_agent,
certs,
agent: config.new_agent(),
auth,
retry,
allow_invalid_certs,
rate_limit: RateLimitState::new(),
log_level: AtomicU8::new(LEVEL_DEFAULT),
requests_gate: OnceLock::new(),
@@ -184,6 +172,10 @@ 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
}
@@ -394,22 +386,17 @@ impl HttpClient {
let result = if let Some(payload) = body {
let mut req = self
.inner
.agent_for(url)
.agent
.post(url)
.header("Authorization", auth)
.header("Accept", "application/json");
if let Some(ct) = content_type {
req = req.header("Content-Type", ct);
}
// 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)
req.send(payload)
} else {
self.inner
.agent_for(url)
.agent
.get(url)
.header("Authorization", auth)
.header("Accept", "application/json")
@@ -657,66 +644,6 @@ 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() {
@@ -784,7 +711,7 @@ mod tests {
token: "t".to_owned(),
},
RetryPolicy::new(0),
CertOverride::none(),
false,
);
let body = br#"{"type":"urn:ietf:params:jmap:error:limit","limit":"someServerLimit"}"#;
assert!(matches!(
@@ -800,7 +727,7 @@ mod tests {
token: "t".to_owned(),
},
RetryPolicy::new(0),
CertOverride::none(),
false,
);
let body = br#"{"type":"urn:ietf:params:jmap:error:limit","limit":"maxSizeRequest"}"#;
assert!(matches!(
@@ -816,7 +743,7 @@ mod tests {
token: "t".to_owned(),
},
RetryPolicy::new(0),
CertOverride::none(),
false,
);
let body =
br#"{"type":"urn:ietf:params:jmap:error:limit","limit":"maxConcurrentRequests"}"#;
@@ -845,7 +772,7 @@ mod tests {
token: "t".to_owned(),
},
RetryPolicy::new(0),
CertOverride::none(),
false,
);
client.set_limits(&limits_with(10, 4, 4));
let err = client
@@ -869,7 +796,7 @@ mod tests {
token: "t".to_owned(),
},
RetryPolicy::new(0),
CertOverride::none(),
false,
);
client.set_limits(&limits_with(1024, 4, 4));
let err = client
@@ -921,7 +848,7 @@ mod tests {
token: "t".to_owned(),
},
RetryPolicy::new(0),
CertOverride::none(),
false,
);
assert_eq!(client.retries_observed(), 0);
assert_eq!(client.retry_after_sleeps(), 0);
+1 -2
View File
@@ -786,7 +786,6 @@ fn decode_set(mr: &MethodCall) -> SetOutcome {
#[cfg(test)]
mod tests {
use crate::net::CertOverride;
use std::cell::Cell;
use super::*;
@@ -799,7 +798,7 @@ mod tests {
token: "t".to_owned(),
},
RetryPolicy::new(max_retries),
CertOverride::none(),
false,
)
}
-2
View File
@@ -1,6 +1,5 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -17,7 +16,6 @@ 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
@@ -1,284 +0,0 @@
/*
* 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));
}
}
-30
View File
@@ -106,34 +106,6 @@ 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);
}
@@ -472,8 +444,6 @@ mod keyed;
mod sieve;
mod sieve_names;
mod uidtype;
mod email;
+29 -114
View File
@@ -10,13 +10,12 @@ use std::collections::{HashMap, HashSet};
use serde_json::{Value, json};
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::db;
use crate::error::Error;
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::jmap::request::Request;
use crate::logging::Logger;
use crate::sync::import_jmap::mapping::{SIEVE_SELECT, row_to_sieve_script};
use crate::sync::{Context, TypeCounts};
use crate::types::ObjectType;
@@ -65,67 +64,32 @@ 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 {
// 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)) {
match content_differs(ctx, net, *blob_local, 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;
}
Ok(true) => match uploader.upload_with(*blob_local, "application/sieve") {
Ok(blob) => updates.push((id.clone(), json!({ "blobId": blob.0 }))),
Err(e) => {
logger.warn(&format!("SieveScript {id}: upload for update failed: {e}"));
counts.failed += 1;
}
}
},
Err(e) => {
logger.warn(&format!("SieveScript {label}: not compared: {e}"));
logger.warn(&format!("SieveScript {id}: 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 = 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 blob_id = 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()));
@@ -149,7 +113,7 @@ pub fn reconcile(
}
None => {
for (cid, err) in &outcome.not_created {
logger.warn(&format!("SieveScript {label} ({cid}) not created: {err}"));
logger.warn(&format!("SieveScript {cid} not created: {err}"));
}
counts.failed += 1;
continue;
@@ -166,15 +130,6 @@ pub fn reconcile(
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();
@@ -186,20 +141,8 @@ pub fn reconcile(
json!({ "accountId": net.account })
};
req.call("SieveScript/set", args, "a");
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}")),
}
if let Err(e) = req.send(&net.client, &net.api) {
logger.warn(&format!("SieveScript activation failed: {e}"));
}
}
@@ -217,48 +160,20 @@ pub fn reconcile(
})
}
/// 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> {
/// Whether the target's copy of a script differs from the archive's. A target
/// that reports no blob is taken as different, so the archive's is written.
fn content_differs(
ctx: &Context,
net: &Net,
blob_local: i64,
target_blob: Option<&String>,
) -> Result<bool, Error> {
let Some(target_blob) = target_blob else {
return Ok(true);
};
let ours = db::blobs::blob_bytes(&ctx.conn, blob_local)
.map_err(|e| Error::Partial(e.to_string()))?
.ok_or_else(|| Error::Partial(format!("blob local id {blob_local} missing")))?;
let theirs = blobxfer::download_bytes(
&net.client,
&net.session,
@@ -268,5 +183,5 @@ fn content_differs(net: &Net, ours: &[u8], target_blob: Option<&String>) -> Resu
"script.sieve",
)
.map_err(Error::from)?;
Ok(ours != theirs.as_slice())
Ok(ours != theirs)
}
-342
View File
@@ -1,342 +0,0 @@
/*
* 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(&[]));
}
}
+1 -17
View File
@@ -187,19 +187,13 @@ fn build_wire(
}
}
/// 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);
let stamp = |v: Option<&Value>| v.and_then(Value::as_str).map(str::to_owned);
if let (Some(ours), Some(theirs)) = (stamp(wire.get("updated")), stamp(target.get("updated")))
&& ours <= theirs
{
@@ -241,16 +235,6 @@ mod tests {
);
}
#[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"});
+1 -3
View File
@@ -1,6 +1,5 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -19,7 +18,6 @@ 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 {
@@ -107,7 +105,7 @@ fn run_into(
let client = DavClient::new(
config.auth.to_jmap_auth(),
RetryPolicy::new(common.max_retries),
CertOverride::for_url(common.allow_invalid_certs, &config.url),
common.allow_invalid_certs,
);
client.set_logger(logger);
+38 -58
View File
@@ -20,7 +20,6 @@ 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 {
@@ -46,8 +45,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)?;
let (discovery, certs) = run_autodiscover(&config, &acquired, common.allow_invalid_certs)?;
let (auth, acquired) = resolve_auth(&config.auth, common.allow_invalid_certs)?;
let discovery = run_autodiscover(&config, &acquired, common.allow_invalid_certs)?;
if logger.enabled(LEVEL_PROGRESS) {
eprintln!(
"EWS discovery: url={} source={:?}",
@@ -78,7 +77,7 @@ pub fn run(common: CommonConfig, config: EwsImportConfig) -> Result<Summary, Err
let client = EwsClient::new(
auth,
RetryPolicy::new(common.max_retries),
certs.narrowed_to(&discovery.ews_url),
common.allow_invalid_certs,
);
client.set_logger(logger);
if matches!(config.mailbox_kind, MailboxKind::PublicFolders) {
@@ -89,7 +88,13 @@ 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, logger);
spawn_token_refresher(
&client,
&config.auth,
&acquired,
common.allow_invalid_certs,
logger,
);
let username = match &config.auth {
EwsAuth::Basic { user, .. } => user.clone(),
@@ -206,7 +211,10 @@ fn run_dry(
Ok(summary)
}
fn resolve_auth(auth: &EwsAuth) -> Result<(Auth, Option<AcquiredToken>), Error> {
fn resolve_auth(
auth: &EwsAuth,
allow_invalid_certs: bool,
) -> Result<(Auth, Option<AcquiredToken>), Error> {
match auth {
EwsAuth::Basic { user, password } => Ok((
Auth::Basic {
@@ -216,9 +224,12 @@ fn resolve_auth(auth: &EwsAuth) -> Result<(Auth, Option<AcquiredToken>), Error>
None,
)),
EwsAuth::Bearer { token } => {
let acq = acquire(&OAuthFlow::PreAcquired {
token: token.clone(),
})
let acq = acquire(
&OAuthFlow::PreAcquired {
token: token.clone(),
},
allow_invalid_certs,
)
.map_err(Error::from)?;
Ok((
Auth::Bearer {
@@ -228,7 +239,7 @@ fn resolve_auth(auth: &EwsAuth) -> Result<(Auth, Option<AcquiredToken>), Error>
))
}
EwsAuth::OAuth(flow) => {
let acq = acquire(flow).map_err(Error::from)?;
let acq = acquire(flow, allow_invalid_certs).map_err(Error::from)?;
Ok((
Auth::Bearer {
token: acq.access_token.clone(),
@@ -243,36 +254,19 @@ fn run_autodiscover(
config: &EwsImportConfig,
acquired: &Option<AcquiredToken>,
allow_invalid_certs: bool,
) -> Result<(DiscoveryResult, CertOverride), Error> {
) -> Result<DiscoveryResult, Error> {
let email = config
.mailbox
.clone()
.or_else(|| acquired.as_ref().and_then(|a| a.upn.clone()));
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(),
}
let result = discover(
config.url.as_deref(),
email.as_deref(),
None,
allow_invalid_certs,
)
.map_err(Error::from)?;
Ok(result)
}
fn resolve_mailbox(
@@ -376,6 +370,7 @@ fn spawn_token_refresher(
client: &EwsClient,
auth: &EwsAuth,
initial: &Option<AcquiredToken>,
allow_invalid_certs: bool,
logger: crate::logging::Logger,
) {
let flow = match auth {
@@ -407,9 +402,14 @@ 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)
crate::exchange_ews::oauth::refresh_with_token(
tenant,
client_id,
rt,
allow_invalid_certs,
)
} else {
crate::exchange_ews::oauth::acquire(&flow)
crate::exchange_ews::oauth::acquire(&flow, allow_invalid_certs)
};
match result {
Ok(tok) => {
@@ -451,26 +451,6 @@ 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!(
+18 -7
View File
@@ -20,7 +20,6 @@ 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)]
@@ -69,11 +68,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)?;
let acquired = acquire_with_flow(&config.auth, common.allow_invalid_certs)?;
let client = GraphClient::new(
acquired.access_token.clone(),
RetryPolicy::new(common.max_retries),
CertOverride::for_url(common.allow_invalid_certs, &config.api_base),
common.allow_invalid_certs,
);
client.set_logger(logger);
@@ -122,7 +121,13 @@ pub fn run(common: CommonConfig, config: GraphImportConfig) -> Result<Summary, E
&principal.user_principal_name,
)?;
let _refresher = spawn_token_refresher(&client, &config.auth, &acquired, logger);
let _refresher = spawn_token_refresher(
&client,
&config.auth,
&acquired,
common.allow_invalid_certs,
logger,
);
let mut summary = Summary::default();
let mut mailbox_counts = TypeCounts::default();
@@ -277,7 +282,7 @@ pub fn enumerate_mail_folders(
Ok(all)
}
fn acquire_with_flow(auth: &GraphAuth) -> Result<AcquiredToken, Error> {
fn acquire_with_flow(auth: &GraphAuth, allow_invalid_certs: bool) -> Result<AcquiredToken, Error> {
let flow = match auth {
GraphAuth::PreAcquired { token } => OAuthFlow::PreAcquired {
token: token.clone(),
@@ -290,7 +295,7 @@ fn acquire_with_flow(auth: &GraphAuth) -> Result<AcquiredToken, Error> {
client_id: client_id.clone(),
},
};
acquire(&flow).map_err(Error::from)
acquire(&flow, allow_invalid_certs).map_err(Error::from)
}
fn resolve_endpoints(config: &GraphImportConfig, client: &GraphClient) -> Result<Endpoints, Error> {
@@ -374,6 +379,7 @@ 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 {
@@ -410,7 +416,12 @@ fn spawn_token_refresher(
break;
}
}
match refresh_access_token(&authority, &client_id, &refresh_token) {
match refresh_access_token(
&authority,
&client_id,
&refresh_token,
allow_invalid_certs,
) {
Ok(tok) => {
client.set_bearer(tok.access_token.clone());
if let Some(new_refresh) = tok.refresh_token {
+1 -3
View File
@@ -1,6 +1,5 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* SPDX-License-Identifier: Apache-2.0 OR MIT
*/
@@ -42,7 +41,6 @@ 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 {
@@ -136,7 +134,7 @@ impl Context {
let client = HttpClient::new(
connect.auth.clone(),
RetryPolicy::new(common.max_retries),
CertOverride::for_url(common.allow_invalid_certs, &connect.url),
common.allow_invalid_certs,
);
Ok(Context {
conn,
+1 -2
View File
@@ -11,7 +11,6 @@ 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 {
@@ -21,7 +20,7 @@ fn admin_client() -> HttpClient {
password: seeder::ADMIN_PASSWORD.into(),
},
RetryPolicy::new(5),
CertOverride::for_url(true, shared_stalwart().base_url()),
true,
)
}
+1 -2
View File
@@ -13,7 +13,6 @@ 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(
@@ -22,7 +21,7 @@ fn client(retries: u32) -> DavClient {
password: "p".into(),
},
RetryPolicy::new(retries),
CertOverride::none(),
false,
)
}
+2 -3
View File
@@ -22,7 +22,6 @@ 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";
@@ -34,7 +33,7 @@ fn client(retries: u32) -> EwsClient {
token: "t".to_owned(),
},
RetryPolicy::new(retries),
CertOverride::none(),
false,
)
}
@@ -64,7 +63,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, &CertOverride::none()).unwrap();
let r = discover(Some(url), None, None, false).unwrap();
assert_eq!(r.source, DiscoverySource::SuppliedUrl);
assert_eq!(r.ews_url, url);
}
+2 -11
View File
@@ -21,7 +21,6 @@ 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;
@@ -29,11 +28,7 @@ static INIT: Once = Once::new();
fn client_with_retries(retries: u32) -> GraphClient {
INIT.call_once(|| {});
GraphClient::new(
"BEARER".to_owned(),
RetryPolicy::new(retries),
CertOverride::none(),
)
GraphClient::new("BEARER".to_owned(), RetryPolicy::new(retries), false)
}
fn url_message_collection(server_url: &str, folder: &str, top: usize) -> String {
@@ -1686,11 +1681,7 @@ 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),
CertOverride::none(),
);
let client = GraphClient::new("EXPIRED".to_owned(), RetryPolicy::new(0), false);
let url = format!("{base}/me");
let err = client.get(&url, Accept::Json).unwrap_err();
assert!(matches!(err, GraphError::Auth(_)));
+1 -2
View File
@@ -17,7 +17,6 @@ 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 {
@@ -27,7 +26,7 @@ fn client(retries: u32) -> HttpClient {
password: "p".into(),
},
RetryPolicy::new(retries),
CertOverride::none(),
false,
)
}
-309
View File
@@ -2788,205 +2788,6 @@ fn export_sieve_script_matched_by_name_with_the_same_content_is_left_alone() {
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();
@@ -5190,113 +4991,3 @@ fn export_rerun_leaves_a_contact_alone_when_the_target_is_newer() {
);
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);
}
+5 -22
View File
@@ -16,7 +16,6 @@ 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;
@@ -786,11 +785,7 @@ 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),
CertOverride::for_url(true, base_url()),
);
let client = HttpClient::new(basic("test6"), RetryPolicy::new(5), true);
let session = Session::discover(&client, base_url()).expect("discover target session");
let api = session.api_url.clone();
@@ -1052,11 +1047,7 @@ 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),
CertOverride::for_url(true, base_url()),
);
let client = HttpClient::new(basic("test1"), RetryPolicy::new(20), true);
let session = Session::discover(&client, base_url()).expect("discover session");
let server_limits = session.core_limits().expect("core limits");
@@ -1157,7 +1148,7 @@ impl JmapSettingsGuard {
password: seeder::ADMIN_PASSWORD.to_owned(),
},
RetryPolicy::new(5),
CertOverride::for_url(true, base_url()),
true,
);
let session = Session::discover(&admin, base_url()).expect("admin discover");
let admin_account = session
@@ -1278,11 +1269,7 @@ 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),
CertOverride::for_url(true, base_url()),
);
let client = HttpClient::new(basic("test1"), RetryPolicy::new(20), true);
let session = Session::discover(&client, base_url()).expect("discover session");
let limits = session.core_limits().expect("core limits");
client.set_limits(&limits);
@@ -1359,11 +1346,7 @@ 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),
CertOverride::for_url(true, base_url()),
);
let client = HttpClient::new(basic("test1"), RetryPolicy::new(5), true);
let session = Session::discover(&client, base_url()).expect("session discovered");
let account = account::resolve(
&AccountSelector::Id(acc.account_id.clone()),