+1
-1
@@ -101,7 +101,7 @@ struct GlobalArgs {
|
||||
long_help = "Increase log verbosity (repeatable, max -vvv).\n \
|
||||
(default) per-type start/finish lines and a final summary\n \
|
||||
-v per-chunk progress, account resolution, retry notices\n \
|
||||
-vv every JMAP call (method, accountId, args, status, timing)\n \
|
||||
-vv every protocol call (method/command, target, status, timing)\n \
|
||||
-vvv full request/response bodies and backoff delays (stderr)"
|
||||
)]
|
||||
verbose: u8,
|
||||
|
||||
+92
-13
@@ -7,7 +7,7 @@
|
||||
use std::io::{BufReader, Read};
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicU8, AtomicU64, Ordering};
|
||||
use std::time::Duration;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
use ureq::Agent;
|
||||
use ureq::Body;
|
||||
@@ -20,7 +20,7 @@ use crate::dav::retry::{DavOutcome, classify};
|
||||
use crate::jmap::error::JmapError;
|
||||
use crate::jmap::http::{Auth, RetryPolicy, retry_after_header};
|
||||
use crate::jmap::retry::{self, RateLimitState};
|
||||
use crate::logging::{LEVEL_BODIES, LEVEL_DEFAULT, LEVEL_PROGRESS, Logger};
|
||||
use crate::logging::{HttpCall, LEVEL_BODIES, LEVEL_DEFAULT, LEVEL_PROGRESS, Logger};
|
||||
|
||||
const MAX_BODY: u64 = 512 * 1024 * 1024;
|
||||
const LONG_RETRY_THRESHOLD: Duration = Duration::from_secs(10);
|
||||
@@ -422,6 +422,7 @@ impl DavClient {
|
||||
let mut include_auth = true;
|
||||
loop {
|
||||
self.inner.rate_limit.cooldown().wait();
|
||||
let started = Instant::now();
|
||||
let send = self.build_and_send(WireRequest {
|
||||
method,
|
||||
url: ¤t_url,
|
||||
@@ -473,6 +474,18 @@ impl DavClient {
|
||||
}
|
||||
if (200..300).contains(&status) {
|
||||
self.inner.rate_limit.on_success();
|
||||
logger.trace_http(&HttpCall {
|
||||
proto: "DAV",
|
||||
method,
|
||||
url: ¤t_url,
|
||||
status,
|
||||
elapsed: started.elapsed(),
|
||||
note: Some("streamed multistatus"),
|
||||
request: body,
|
||||
request_type: Some("application/xml"),
|
||||
response: b"",
|
||||
response_type: None,
|
||||
});
|
||||
let parse_base = if current_url == url {
|
||||
base_url
|
||||
} else {
|
||||
@@ -498,6 +511,18 @@ impl DavClient {
|
||||
)));
|
||||
}
|
||||
};
|
||||
logger.trace_http(&HttpCall {
|
||||
proto: "DAV",
|
||||
method,
|
||||
url: ¤t_url,
|
||||
status,
|
||||
elapsed: started.elapsed(),
|
||||
note: None,
|
||||
request: body,
|
||||
request_type: Some("application/xml"),
|
||||
response: &body_bytes,
|
||||
response_type: None,
|
||||
});
|
||||
let disposition = classify(status, &body_bytes);
|
||||
match disposition {
|
||||
DavOutcome::Vanished => {
|
||||
@@ -632,6 +657,12 @@ impl DavClient {
|
||||
}
|
||||
|
||||
fn one_attempt(&self, req: WireRequest<'_>) -> Attempt {
|
||||
let logger = self.logger();
|
||||
let method = req.method;
|
||||
let url = req.url;
|
||||
let request = req.body;
|
||||
let request_type = req.content_type;
|
||||
let started = Instant::now();
|
||||
match self.build_and_send(req) {
|
||||
Ok(mut resp) => {
|
||||
let status = resp.status().as_u16();
|
||||
@@ -645,6 +676,18 @@ impl DavClient {
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.and_then(retry_after_header);
|
||||
if is_redirect_status(status) && location.is_some() {
|
||||
logger.trace_http(&HttpCall {
|
||||
proto: "DAV",
|
||||
method,
|
||||
url,
|
||||
status,
|
||||
elapsed: started.elapsed(),
|
||||
note: location.as_deref(),
|
||||
request,
|
||||
request_type,
|
||||
response: b"",
|
||||
response_type: None,
|
||||
});
|
||||
return Attempt::Ok {
|
||||
status,
|
||||
body: Vec::new(),
|
||||
@@ -666,21 +709,39 @@ impl DavClient {
|
||||
.limit(body_limit)
|
||||
.read_to_vec()
|
||||
{
|
||||
Ok(bytes) => Attempt::Ok {
|
||||
status,
|
||||
body: bytes,
|
||||
retry_after,
|
||||
etag,
|
||||
content_type: content_type_header,
|
||||
last_modified,
|
||||
location,
|
||||
},
|
||||
Ok(bytes) => {
|
||||
logger.trace_http(&HttpCall {
|
||||
proto: "DAV",
|
||||
method,
|
||||
url,
|
||||
status,
|
||||
elapsed: started.elapsed(),
|
||||
note: None,
|
||||
request,
|
||||
request_type,
|
||||
response: &bytes,
|
||||
response_type: content_type_header.as_deref(),
|
||||
});
|
||||
Attempt::Ok {
|
||||
status,
|
||||
body: bytes,
|
||||
retry_after,
|
||||
etag,
|
||||
content_type: content_type_header,
|
||||
last_modified,
|
||||
location,
|
||||
}
|
||||
}
|
||||
Err(e) => Attempt::Transport(JmapError::Transport(format!(
|
||||
"reading response body: {e}"
|
||||
))),
|
||||
}
|
||||
}
|
||||
Err(e) => Attempt::Transport(map_ureq_error(e)),
|
||||
Err(e) => {
|
||||
let err = map_ureq_error(e);
|
||||
logger.trace_http_error("DAV", method, url, &err.to_string(), started.elapsed());
|
||||
Attempt::Transport(err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -694,6 +755,8 @@ impl DavClient {
|
||||
accept: ACCEPT_BINARY,
|
||||
include_auth,
|
||||
};
|
||||
let logger = self.logger();
|
||||
let started = Instant::now();
|
||||
match self.build_and_send(req) {
|
||||
Ok(resp) => {
|
||||
let status = resp.status().as_u16();
|
||||
@@ -706,6 +769,18 @@ impl DavClient {
|
||||
.get("retry-after")
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.and_then(retry_after_header);
|
||||
logger.trace_http(&HttpCall {
|
||||
proto: "DAV",
|
||||
method: "GET",
|
||||
url,
|
||||
status,
|
||||
elapsed: started.elapsed(),
|
||||
note: Some("streamed download"),
|
||||
request: None,
|
||||
request_type: None,
|
||||
response: b"",
|
||||
response_type: None,
|
||||
});
|
||||
let body_reader = resp.into_body().into_reader();
|
||||
AttemptStream::Ok {
|
||||
status,
|
||||
@@ -717,7 +792,11 @@ impl DavClient {
|
||||
retry_after,
|
||||
}
|
||||
}
|
||||
Err(e) => AttemptStream::Transport(map_ureq_error(e)),
|
||||
Err(e) => {
|
||||
let err = map_ureq_error(e);
|
||||
logger.trace_http_error("DAV", "GET", url, &err.to_string(), started.elapsed());
|
||||
AttemptStream::Transport(err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -7,7 +7,7 @@
|
||||
use std::sync::Arc;
|
||||
use std::sync::Mutex;
|
||||
use std::sync::atomic::{AtomicU8, AtomicU64, Ordering};
|
||||
use std::time::Duration;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
use base64::Engine;
|
||||
use base64::engine::general_purpose::STANDARD;
|
||||
@@ -22,7 +22,7 @@ use crate::exchange_ews::soap::{EnvelopeOptions, soap_action, wrap_envelope};
|
||||
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::{LEVEL_BODIES, LEVEL_DEFAULT, LEVEL_PROGRESS, Logger};
|
||||
use crate::logging::{HttpCall, LEVEL_BODIES, LEVEL_DEFAULT, LEVEL_PROGRESS, Logger};
|
||||
|
||||
const MAX_BODY: u64 = 2 * 1024 * 1024 * 1024;
|
||||
const LONG_RETRY_THRESHOLD: Duration = Duration::from_secs(10);
|
||||
@@ -427,10 +427,17 @@ impl EwsClient {
|
||||
if let Some(cookie) = self.affinity_cookie() {
|
||||
req = req.header("X-BackEndOverrideCookie", cookie);
|
||||
}
|
||||
let logger = self.logger();
|
||||
let started = Instant::now();
|
||||
let result = req.send(body.as_bytes());
|
||||
match result {
|
||||
Ok(mut resp) => {
|
||||
let status = resp.status().as_u16();
|
||||
let response_type = resp
|
||||
.headers()
|
||||
.get("content-type")
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.map(str::to_owned);
|
||||
let retry_after = resp
|
||||
.headers()
|
||||
.get("retry-after")
|
||||
@@ -442,17 +449,35 @@ impl EwsClient {
|
||||
*g = Some(cookie);
|
||||
}
|
||||
match resp.body_mut().with_config().limit(MAX_BODY).read_to_vec() {
|
||||
Ok(bytes) => AttemptOutcome::Ok {
|
||||
status,
|
||||
body: bytes,
|
||||
retry_after,
|
||||
},
|
||||
Ok(bytes) => {
|
||||
logger.trace_http(&HttpCall {
|
||||
proto: "EWS",
|
||||
method: action,
|
||||
url,
|
||||
status,
|
||||
elapsed: started.elapsed(),
|
||||
note: None,
|
||||
request: Some(body.as_bytes()),
|
||||
request_type: Some("text/xml"),
|
||||
response: &bytes,
|
||||
response_type: response_type.as_deref(),
|
||||
});
|
||||
AttemptOutcome::Ok {
|
||||
status,
|
||||
body: bytes,
|
||||
retry_after,
|
||||
}
|
||||
}
|
||||
Err(e) => AttemptOutcome::Transport(EwsError::Transport(format!(
|
||||
"reading response body: {e}"
|
||||
))),
|
||||
}
|
||||
}
|
||||
Err(e) => AttemptOutcome::Transport(map_ureq_error(e)),
|
||||
Err(e) => {
|
||||
let err = map_ureq_error(e);
|
||||
logger.trace_http_error("EWS", action, url, &err.to_string(), started.elapsed());
|
||||
AttemptOutcome::Transport(err)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -7,18 +7,19 @@
|
||||
use std::sync::Arc;
|
||||
use std::sync::Mutex;
|
||||
use std::sync::atomic::{AtomicU8, AtomicU64, Ordering};
|
||||
use std::time::Duration;
|
||||
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;
|
||||
use crate::exchange_graph::retry::{HttpClass, classify_http_status, is_throttled};
|
||||
use crate::jmap::http::{RetryPolicy, retry_after_header};
|
||||
use crate::jmap::http::{RetryPolicy, cross_host, retry_after_header};
|
||||
use crate::jmap::retry::{self, RateLimitState};
|
||||
use crate::logging::{LEVEL_BODIES, LEVEL_DEFAULT, LEVEL_PROGRESS, Logger};
|
||||
use crate::logging::{HttpCall, LEVEL_BODIES, LEVEL_DEFAULT, LEVEL_PROGRESS, Logger};
|
||||
|
||||
const MAX_BODY: u64 = 256 * 1024 * 1024;
|
||||
const LONG_RETRY_THRESHOLD: Duration = Duration::from_secs(10);
|
||||
@@ -321,9 +322,12 @@ impl GraphClient {
|
||||
for value in extra_prefer {
|
||||
req = req.header("Prefer", *value);
|
||||
}
|
||||
let logger = self.logger();
|
||||
let started = Instant::now();
|
||||
match req.call() {
|
||||
Ok(mut resp) => {
|
||||
let status = resp.status().as_u16();
|
||||
let final_uri = resp.get_uri().clone();
|
||||
let retry_after = resp
|
||||
.headers()
|
||||
.get("retry-after")
|
||||
@@ -335,18 +339,37 @@ impl GraphClient {
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.map(str::to_owned);
|
||||
match resp.body_mut().with_config().limit(MAX_BODY).read_to_vec() {
|
||||
Ok(bytes) => Attempt::Ok {
|
||||
status,
|
||||
body: bytes,
|
||||
retry_after,
|
||||
content_type,
|
||||
},
|
||||
Ok(bytes) => {
|
||||
warn_on_redirect(&logger, method, url, &final_uri);
|
||||
logger.trace_http(&HttpCall {
|
||||
proto: "Graph",
|
||||
method,
|
||||
url,
|
||||
status,
|
||||
elapsed: started.elapsed(),
|
||||
note: None,
|
||||
request: None,
|
||||
request_type: None,
|
||||
response: &bytes,
|
||||
response_type: content_type.as_deref(),
|
||||
});
|
||||
Attempt::Ok {
|
||||
status,
|
||||
body: bytes,
|
||||
retry_after,
|
||||
content_type,
|
||||
}
|
||||
}
|
||||
Err(e) => Attempt::Transport(GraphError::Transport(format!(
|
||||
"reading response body: {e}"
|
||||
))),
|
||||
}
|
||||
}
|
||||
Err(e) => Attempt::Transport(map_ureq_error(e)),
|
||||
Err(e) => {
|
||||
let err = map_ureq_error(e);
|
||||
logger.trace_http_error("Graph", method, url, &err.to_string(), started.elapsed());
|
||||
Attempt::Transport(err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -387,6 +410,15 @@ struct RetryLog<'a> {
|
||||
body: &'a [u8],
|
||||
}
|
||||
|
||||
fn warn_on_redirect(logger: &Logger, method: &str, requested: &str, final_uri: &Uri) {
|
||||
let landed = final_uri.to_string();
|
||||
if cross_host(requested, &landed) {
|
||||
logger.warn(&format!(
|
||||
"{method} {requested} was redirected across hosts to {landed}; verify the Graph endpoint host is reachable directly without redirection"
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
fn map_ureq_error(err: ureq::Error) -> GraphError {
|
||||
match err {
|
||||
ureq::Error::Io(e) => GraphError::Transport(format!("io: {e}")),
|
||||
|
||||
@@ -6,6 +6,7 @@
|
||||
|
||||
use std::collections::BTreeSet;
|
||||
use std::io::{BufReader, Read, Write};
|
||||
use std::time::Instant;
|
||||
|
||||
use base64::Engine;
|
||||
use base64::engine::general_purpose::STANDARD as BASE64;
|
||||
@@ -14,6 +15,7 @@ use super::command::{self, CommandBuilder};
|
||||
use super::error::{ImapError, NoError};
|
||||
use super::response::{Response, Status, StatusLine, Untagged, parse_response};
|
||||
use super::transport::{Connector, ImapStream};
|
||||
use crate::logging::Logger;
|
||||
|
||||
pub struct ImapClient {
|
||||
reader: BufReader<Box<dyn ImapStream>>,
|
||||
@@ -22,6 +24,7 @@ pub struct ImapClient {
|
||||
pub host: String,
|
||||
closed: bool,
|
||||
utf8_accept: bool,
|
||||
logger: Logger,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
@@ -44,6 +47,7 @@ impl ImapClient {
|
||||
host: &str,
|
||||
port: u16,
|
||||
mode: ConnectMode,
|
||||
logger: Logger,
|
||||
) -> Result<ImapClient, ImapError> {
|
||||
let stream: Box<dyn ImapStream> = match mode {
|
||||
ConnectMode::ImplicitTls => connector.connect_tls(host, port)?,
|
||||
@@ -56,6 +60,7 @@ impl ImapClient {
|
||||
host: host.to_owned(),
|
||||
closed: false,
|
||||
utf8_accept: false,
|
||||
logger,
|
||||
};
|
||||
client.read_greeting()?;
|
||||
client.send_capability()?;
|
||||
@@ -337,12 +342,19 @@ impl ImapClient {
|
||||
.into(),
|
||||
));
|
||||
}
|
||||
let started = Instant::now();
|
||||
self.write_all(&bytes)?;
|
||||
loop {
|
||||
let resp = parse_response(&mut self.reader)?;
|
||||
match resp {
|
||||
Response::Tagged { tag: t, line } if t == tag => {
|
||||
self.update_capabilities_from_code(&line);
|
||||
self.logger.trace_cmd(
|
||||
"IMAP",
|
||||
command,
|
||||
&format!("{} {}", status_word(&line.status), line.text),
|
||||
started.elapsed(),
|
||||
);
|
||||
return match line.status {
|
||||
Status::Ok => Ok(line),
|
||||
Status::No => Err(ImapError::No(line.into_no_error())),
|
||||
@@ -388,6 +400,7 @@ impl ImapClient {
|
||||
.into(),
|
||||
));
|
||||
}
|
||||
let started = Instant::now();
|
||||
self.write_all(&bytes)?;
|
||||
let mut untagged = Vec::new();
|
||||
loop {
|
||||
@@ -395,6 +408,12 @@ impl ImapClient {
|
||||
match resp {
|
||||
Response::Tagged { tag: t, line } if t == tag => {
|
||||
self.update_capabilities_from_code(&line);
|
||||
self.logger.trace_cmd(
|
||||
"IMAP",
|
||||
command,
|
||||
&format!("{} {}", status_word(&line.status), line.text),
|
||||
started.elapsed(),
|
||||
);
|
||||
return match line.status {
|
||||
Status::Ok => Ok(CollectedResponse {
|
||||
tag,
|
||||
@@ -475,6 +494,16 @@ pub struct CollectedResponse {
|
||||
pub untagged: Vec<Untagged>,
|
||||
}
|
||||
|
||||
fn status_word(status: &Status) -> &'static str {
|
||||
match status {
|
||||
Status::Ok => "OK",
|
||||
Status::No => "NO",
|
||||
Status::Bad => "BAD",
|
||||
Status::Bye => "BYE",
|
||||
Status::PreAuth => "PREAUTH",
|
||||
}
|
||||
}
|
||||
|
||||
fn parse_caps(s: &str) -> BTreeSet<String> {
|
||||
s.split_ascii_whitespace().map(|w| w.to_owned()).collect()
|
||||
}
|
||||
@@ -569,6 +598,7 @@ mod tests {
|
||||
host: "test.example".to_owned(),
|
||||
closed: false,
|
||||
utf8_accept: false,
|
||||
logger: Logger::from_flags(true, 0),
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+89
-11
@@ -7,7 +7,7 @@
|
||||
use std::sync::Arc;
|
||||
use std::sync::OnceLock;
|
||||
use std::sync::atomic::{AtomicU8, AtomicU64, Ordering};
|
||||
use std::time::{Duration, SystemTime};
|
||||
use std::time::{Duration, Instant, SystemTime};
|
||||
|
||||
use base64::Engine;
|
||||
use base64::engine::general_purpose::STANDARD;
|
||||
@@ -15,12 +15,13 @@ 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;
|
||||
use crate::jmap::inflight::{Permit, Semaphore};
|
||||
use crate::jmap::retry::{self, Disposition, RateLimitState};
|
||||
use crate::jmap::session::Limits;
|
||||
use crate::logging::{LEVEL_BODIES, LEVEL_DEFAULT, LEVEL_PROGRESS, Logger};
|
||||
use crate::logging::{HttpCall, LEVEL_BODIES, LEVEL_DEFAULT, LEVEL_PROGRESS, Logger};
|
||||
|
||||
const MAX_BODY: u64 = 512 * 1024 * 1024;
|
||||
|
||||
@@ -253,7 +254,7 @@ impl HttpClient {
|
||||
} else {
|
||||
None
|
||||
};
|
||||
self.one_attempt(url, body, content_type)
|
||||
self.one_attempt(method, url, body, content_type)
|
||||
};
|
||||
match attempt_outcome {
|
||||
Attempt::Ok { status, body, .. } if (200..300).contains(&status) => {
|
||||
@@ -366,8 +367,16 @@ impl HttpClient {
|
||||
}
|
||||
}
|
||||
|
||||
fn one_attempt(&self, url: &str, body: Option<&[u8]>, content_type: Option<&str>) -> Attempt {
|
||||
fn one_attempt(
|
||||
&self,
|
||||
method: &str,
|
||||
url: &str,
|
||||
body: Option<&[u8]>,
|
||||
content_type: Option<&str>,
|
||||
) -> Attempt {
|
||||
let auth = self.inner.auth.header_value();
|
||||
let logger = self.logger();
|
||||
let started = Instant::now();
|
||||
let result = if let Some(payload) = body {
|
||||
let mut req = self
|
||||
.inner
|
||||
@@ -390,6 +399,12 @@ impl HttpClient {
|
||||
match result {
|
||||
Ok(mut resp) => {
|
||||
let status = resp.status().as_u16();
|
||||
let final_uri = resp.get_uri().clone();
|
||||
let response_type = resp
|
||||
.headers()
|
||||
.get("content-type")
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.map(str::to_owned);
|
||||
let retry_after = resp
|
||||
.headers()
|
||||
.get("retry-after")
|
||||
@@ -411,18 +426,37 @@ impl HttpClient {
|
||||
})
|
||||
.collect();
|
||||
match resp.body_mut().with_config().limit(MAX_BODY).read_to_vec() {
|
||||
Ok(bytes) => Attempt::Ok {
|
||||
status,
|
||||
body: bytes,
|
||||
retry_after,
|
||||
rate_limit_headers,
|
||||
},
|
||||
Ok(bytes) => {
|
||||
warn_on_redirect(&logger, method, url, &final_uri);
|
||||
logger.trace_http(&HttpCall {
|
||||
proto: "JMAP",
|
||||
method,
|
||||
url,
|
||||
status,
|
||||
elapsed: started.elapsed(),
|
||||
note: None,
|
||||
request: body,
|
||||
request_type: content_type,
|
||||
response: &bytes,
|
||||
response_type: response_type.as_deref(),
|
||||
});
|
||||
Attempt::Ok {
|
||||
status,
|
||||
body: bytes,
|
||||
retry_after,
|
||||
rate_limit_headers,
|
||||
}
|
||||
}
|
||||
Err(e) => Attempt::Transport(JmapError::Transport(format!(
|
||||
"reading response body: {e}"
|
||||
))),
|
||||
}
|
||||
}
|
||||
Err(e) => Attempt::Transport(map_ureq_error(e)),
|
||||
Err(e) => {
|
||||
let err = map_ureq_error(e);
|
||||
logger.trace_http_error("JMAP", method, url, &err.to_string(), started.elapsed());
|
||||
Attempt::Transport(err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -498,6 +532,26 @@ struct RetryLog<'a> {
|
||||
body: &'a [u8],
|
||||
}
|
||||
|
||||
fn warn_on_redirect(logger: &Logger, method: &str, requested: &str, final_uri: &Uri) {
|
||||
let landed = final_uri.to_string();
|
||||
if cross_host(requested, &landed) {
|
||||
logger.warn(&format!(
|
||||
"{method} {requested} was redirected across hosts to {landed}; a \
|
||||
JMAP endpoint must not redirect: verify the server's advertised base \
|
||||
URL and that the request host matches its configured hostname"
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
pub fn cross_host(requested: &str, landed: &str) -> bool {
|
||||
match (url::Url::parse(requested), url::Url::parse(landed)) {
|
||||
(Ok(a), Ok(b)) => {
|
||||
a.host_str() != b.host_str() || a.port_or_known_default() != b.port_or_known_default()
|
||||
}
|
||||
_ => false,
|
||||
}
|
||||
}
|
||||
|
||||
fn transport_disposition(err: &JmapError) -> Disposition {
|
||||
match err {
|
||||
JmapError::Connect(_) => Disposition::Fatal,
|
||||
@@ -600,6 +654,30 @@ mod tests {
|
||||
assert_eq!(auth.header_value(), "Bearer abc.def");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn cross_host_detects_hostname_change() {
|
||||
assert!(cross_host(
|
||||
"https://mail.example.com/jmap",
|
||||
"https://login.example.net/sso"
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn cross_host_ignores_same_host_path_change() {
|
||||
assert!(!cross_host(
|
||||
"https://mail.example.com/jmap",
|
||||
"https://mail.example.com/jmap/api"
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn cross_host_flags_port_change() {
|
||||
assert!(cross_host(
|
||||
"https://mail.example.com/jmap",
|
||||
"https://mail.example.com:8443/jmap"
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn truncate_does_not_panic_on_multibyte_boundary() {
|
||||
let mut body = vec![b'a'; 510];
|
||||
|
||||
+232
@@ -4,12 +4,16 @@
|
||||
* SPDX-License-Identifier: Apache-2.0 OR MIT
|
||||
*/
|
||||
|
||||
use std::time::Duration;
|
||||
|
||||
pub const LEVEL_QUIET: u8 = 0;
|
||||
pub const LEVEL_DEFAULT: u8 = 1;
|
||||
pub const LEVEL_PROGRESS: u8 = 2;
|
||||
pub const LEVEL_METHOD: u8 = 3;
|
||||
pub const LEVEL_BODIES: u8 = 4;
|
||||
|
||||
pub const TRACE_BODY_CAP: usize = 8192;
|
||||
|
||||
#[derive(Debug, Clone, Copy)]
|
||||
pub struct Logger {
|
||||
level: u8,
|
||||
@@ -46,6 +50,147 @@ impl Logger {
|
||||
pub fn error(&self, message: &str) {
|
||||
eprintln!("error: {message}");
|
||||
}
|
||||
|
||||
pub fn trace_http(&self, call: &HttpCall<'_>) {
|
||||
if !self.enabled(LEVEL_METHOD) {
|
||||
return;
|
||||
}
|
||||
match call.note {
|
||||
Some(note) => eprintln!(
|
||||
"{} {} {} -> {} ({} ms) {note}",
|
||||
call.proto,
|
||||
call.method,
|
||||
call.url,
|
||||
call.status,
|
||||
call.elapsed.as_millis()
|
||||
),
|
||||
None => eprintln!(
|
||||
"{} {} {} -> {} ({} ms)",
|
||||
call.proto,
|
||||
call.method,
|
||||
call.url,
|
||||
call.status,
|
||||
call.elapsed.as_millis()
|
||||
),
|
||||
}
|
||||
if self.enabled(LEVEL_BODIES) {
|
||||
if let Some(request) = call.request {
|
||||
trace_body('>', request, call.request_type);
|
||||
}
|
||||
trace_body('<', call.response, call.response_type);
|
||||
}
|
||||
}
|
||||
|
||||
pub fn trace_cmd(&self, proto: &str, command: &str, status: &str, elapsed: Duration) {
|
||||
if !self.enabled(LEVEL_METHOD) {
|
||||
return;
|
||||
}
|
||||
let command = cap(&redact_wire(command), command.len());
|
||||
eprintln!("{proto} {command} -> {status} ({} ms)", elapsed.as_millis());
|
||||
}
|
||||
|
||||
pub fn trace_http_error(&self, proto: &str, method: &str, url: &str, error: &str, elapsed: Duration) {
|
||||
if !self.enabled(LEVEL_METHOD) {
|
||||
return;
|
||||
}
|
||||
eprintln!(
|
||||
"{proto} {method} {url} -> transport error: {error} ({} ms)",
|
||||
elapsed.as_millis()
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
pub struct HttpCall<'a> {
|
||||
pub proto: &'a str,
|
||||
pub method: &'a str,
|
||||
pub url: &'a str,
|
||||
pub status: u16,
|
||||
pub elapsed: Duration,
|
||||
pub note: Option<&'a str>,
|
||||
pub request: Option<&'a [u8]>,
|
||||
pub request_type: Option<&'a str>,
|
||||
pub response: &'a [u8],
|
||||
pub response_type: Option<&'a str>,
|
||||
}
|
||||
|
||||
fn trace_body(marker: char, body: &[u8], content_type: Option<&str>) {
|
||||
for line in body_for_log(body, content_type).lines() {
|
||||
eprintln!(" {marker} {line}");
|
||||
}
|
||||
}
|
||||
|
||||
pub fn body_for_log(body: &[u8], content_type: Option<&str>) -> String {
|
||||
if body.is_empty() {
|
||||
return "<empty body>".to_owned();
|
||||
}
|
||||
let is_json = content_type.is_some_and(|c| c.to_ascii_lowercase().contains("json"));
|
||||
if is_json
|
||||
&& let Ok(value) = serde_json::from_slice::<serde_json::Value>(body)
|
||||
&& let Ok(pretty) = serde_json::to_string_pretty(&value)
|
||||
{
|
||||
return cap(&pretty, body.len());
|
||||
}
|
||||
let head = safe_head(body, TRACE_BODY_CAP);
|
||||
match std::str::from_utf8(head) {
|
||||
Ok(text) if is_texty(text) => {
|
||||
if body.len() > head.len() {
|
||||
format!("{text}\n... <{} bytes total>", body.len())
|
||||
} else {
|
||||
text.to_owned()
|
||||
}
|
||||
}
|
||||
_ => blob_summary(body),
|
||||
}
|
||||
}
|
||||
|
||||
fn safe_head(body: &[u8], max: usize) -> &[u8] {
|
||||
if body.len() <= max {
|
||||
return body;
|
||||
}
|
||||
let mut end = max;
|
||||
while end > 0 && (body[end] & 0xC0) == 0x80 {
|
||||
end -= 1;
|
||||
}
|
||||
&body[..end]
|
||||
}
|
||||
|
||||
fn is_texty(text: &str) -> bool {
|
||||
!text
|
||||
.chars()
|
||||
.any(|c| c.is_control() && !matches!(c, '\n' | '\r' | '\t'))
|
||||
}
|
||||
|
||||
fn blob_summary(body: &[u8]) -> String {
|
||||
format!(
|
||||
"<{} bytes, blake3={}>",
|
||||
body.len(),
|
||||
blake3::hash(body).to_hex()
|
||||
)
|
||||
}
|
||||
|
||||
fn cap(text: &str, total_bytes: usize) -> String {
|
||||
if text.len() <= TRACE_BODY_CAP {
|
||||
return text.to_owned();
|
||||
}
|
||||
let end = text.floor_char_boundary(TRACE_BODY_CAP);
|
||||
format!("{}\n... <{total_bytes} bytes total>", &text[..end])
|
||||
}
|
||||
|
||||
pub fn redact_wire(command: &str) -> String {
|
||||
let trimmed = command.trim_end_matches(['\r', '\n']);
|
||||
let mut parts = trimmed.splitn(3, ' ');
|
||||
let first = parts.next().unwrap_or("");
|
||||
if first.eq_ignore_ascii_case("LOGIN") {
|
||||
return format!("{first} <redacted>");
|
||||
}
|
||||
if first.eq_ignore_ascii_case("AUTHENTICATE") {
|
||||
let mechanism = parts.next().unwrap_or("");
|
||||
return match parts.next() {
|
||||
Some(_) => format!("{first} {mechanism} <redacted>"),
|
||||
None => format!("{first} {mechanism}").trim_end().to_owned(),
|
||||
};
|
||||
}
|
||||
trimmed.to_owned()
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
@@ -74,4 +219,91 @@ mod tests {
|
||||
assert!(l.enabled(LEVEL_PROGRESS));
|
||||
assert!(!l.enabled(LEVEL_METHOD));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn redact_wire_hides_login_arguments() {
|
||||
assert_eq!(redact_wire("LOGIN \"alice\" \"s3cret\""), "LOGIN <redacted>");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn redact_wire_hides_authenticate_initial_response() {
|
||||
assert_eq!(
|
||||
redact_wire("AUTHENTICATE PLAIN AGFsaWNlAHMzY3JldA=="),
|
||||
"AUTHENTICATE PLAIN <redacted>"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn redact_wire_keeps_authenticate_without_initial_response() {
|
||||
assert_eq!(redact_wire("AUTHENTICATE PLAIN"), "AUTHENTICATE PLAIN");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn redact_wire_passes_through_ordinary_commands() {
|
||||
assert_eq!(redact_wire("SELECT INBOX\r\n"), "SELECT INBOX");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn body_for_log_pretty_prints_json() {
|
||||
let out = body_for_log(br#"{"a":1}"#, Some("application/json"));
|
||||
assert!(out.contains("\"a\": 1"), "got {out}");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn body_for_log_summarises_binary_with_blake3() {
|
||||
let out = body_for_log(&[0u8, 159, 146, 150], Some("application/octet-stream"));
|
||||
assert!(out.starts_with("<4 bytes, blake3="), "got {out}");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn body_for_log_reports_empty() {
|
||||
assert_eq!(body_for_log(&[], None), "<empty body>");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn body_for_log_caps_long_text_without_scanning_whole_blob() {
|
||||
let body = vec![b'a'; TRACE_BODY_CAP * 4];
|
||||
let out = body_for_log(&body, Some("text/plain"));
|
||||
assert!(out.contains("bytes total>"), "got tail {}", &out[out.len() - 40..]);
|
||||
assert!(out.len() < body.len());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn body_for_log_handles_multibyte_head_boundary() {
|
||||
let mut body = vec![b'a'; TRACE_BODY_CAP - 1];
|
||||
body.extend_from_slice("\u{1F4A9}".as_bytes());
|
||||
body.extend_from_slice(&[b'b'; 32]);
|
||||
let out = body_for_log(&body, Some("text/plain"));
|
||||
assert!(out.contains("bytes total>"), "got {out}");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn trace_emit_paths_do_not_panic_at_max_verbosity() {
|
||||
let logger = Logger::from_flags(false, 9);
|
||||
logger.trace_http(&HttpCall {
|
||||
proto: "JMAP",
|
||||
method: "POST",
|
||||
url: "https://mail.example.com/jmap",
|
||||
status: 200,
|
||||
elapsed: Duration::from_millis(3),
|
||||
note: Some("Mailbox/get"),
|
||||
request: Some(br#"{"using":["urn:ietf:params:jmap:core"]}"#),
|
||||
request_type: Some("application/json"),
|
||||
response: &[0u8, 200, 1, 159, 146, 150],
|
||||
response_type: Some("application/octet-stream"),
|
||||
});
|
||||
logger.trace_cmd(
|
||||
"IMAP",
|
||||
"LOGIN \"alice\" \"s3cret\"",
|
||||
"OK done",
|
||||
Duration::from_millis(1),
|
||||
);
|
||||
logger.trace_http_error(
|
||||
"Graph",
|
||||
"GET",
|
||||
"https://graph.microsoft.com/v1.0/me",
|
||||
"host not found",
|
||||
Duration::from_millis(9),
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -5,11 +5,13 @@
|
||||
*/
|
||||
|
||||
use std::io::{BufRead, BufReader, Read, Write};
|
||||
use std::time::Instant;
|
||||
|
||||
use base64::Engine;
|
||||
use base64::engine::general_purpose::STANDARD as BASE64;
|
||||
|
||||
use crate::imap::transport::{Connector, ImapStream};
|
||||
use crate::logging::Logger;
|
||||
|
||||
use super::command;
|
||||
use super::error::SieveError;
|
||||
@@ -30,6 +32,7 @@ pub struct SieveClient {
|
||||
pub capabilities: Capabilities,
|
||||
closed: bool,
|
||||
fresh_post_auth_caps: bool,
|
||||
logger: Logger,
|
||||
}
|
||||
|
||||
impl SieveClient {
|
||||
@@ -39,6 +42,7 @@ impl SieveClient {
|
||||
port: u16,
|
||||
mode: ConnectMode,
|
||||
allow_cleartext: bool,
|
||||
logger: Logger,
|
||||
) -> Result<SieveClient, SieveError> {
|
||||
let stream: Box<dyn ImapStream> = match mode {
|
||||
ConnectMode::ImplicitTls => connector
|
||||
@@ -54,6 +58,7 @@ impl SieveClient {
|
||||
capabilities: Capabilities::default(),
|
||||
closed: false,
|
||||
fresh_post_auth_caps: false,
|
||||
logger,
|
||||
};
|
||||
client.read_initial_capabilities()?;
|
||||
if matches!(mode, ConnectMode::StartTls) {
|
||||
@@ -143,14 +148,24 @@ impl SieveClient {
|
||||
}
|
||||
|
||||
pub fn run(&mut self, command: &str) -> Result<ResponseBlock, SieveError> {
|
||||
let started = Instant::now();
|
||||
self.write_all(command.as_bytes())?;
|
||||
let block = read_response(&mut self.reader)?;
|
||||
self.logger
|
||||
.trace_cmd("SIEVE", command, &sieve_status(&block.status), started.elapsed());
|
||||
finish_block(block, &mut self.closed)
|
||||
}
|
||||
|
||||
pub fn run_raw(&mut self, bytes: &[u8]) -> Result<ResponseBlock, SieveError> {
|
||||
let started = Instant::now();
|
||||
self.write_all(bytes)?;
|
||||
let block = read_response(&mut self.reader)?;
|
||||
self.logger.trace_cmd(
|
||||
"SIEVE",
|
||||
&String::from_utf8_lossy(bytes),
|
||||
&sieve_status(&block.status),
|
||||
started.elapsed(),
|
||||
);
|
||||
finish_block(block, &mut self.closed)
|
||||
}
|
||||
|
||||
@@ -300,6 +315,7 @@ impl SieveClient {
|
||||
capabilities: Capabilities::default(),
|
||||
closed: false,
|
||||
fresh_post_auth_caps: false,
|
||||
logger: Logger::from_flags(true, 0),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -345,6 +361,15 @@ fn finish_block(block: ResponseBlock, closed: &mut bool) -> Result<ResponseBlock
|
||||
}
|
||||
}
|
||||
|
||||
fn sieve_status(status: &StatusLine) -> String {
|
||||
let word = match status.status {
|
||||
Status::Ok => "OK",
|
||||
Status::No => "NO",
|
||||
Status::Bye => "BYE",
|
||||
};
|
||||
format!("{word} {}", status.text)
|
||||
}
|
||||
|
||||
fn noop_stream() -> Box<dyn ImapStream> {
|
||||
Box::new(NoopStream)
|
||||
}
|
||||
|
||||
@@ -155,6 +155,7 @@ fn reconnect_control(client: &mut ImapClient, ctx: &ControlCtx) -> Result<(), Im
|
||||
&ctx.endpoint.host,
|
||||
ctx.endpoint.port,
|
||||
ctx.mode,
|
||||
ctx.logger,
|
||||
)?;
|
||||
*client = new_client;
|
||||
authenticate_client(client, &ctx.auth).map_err(|e| ImapError::AuthFailed(e.to_string()))?;
|
||||
@@ -224,7 +225,7 @@ pub fn run(common: CommonConfig, config: ImapImportConfig) -> Result<Summary, Er
|
||||
} else {
|
||||
ConnectMode::StartTls
|
||||
};
|
||||
let mut client = ImapClient::connect(&connector, &endpoint.host, endpoint.port, mode)
|
||||
let mut client = ImapClient::connect(&connector, &endpoint.host, endpoint.port, mode, logger)
|
||||
.map_err(|e| Error::Connection(e.to_string()))?;
|
||||
|
||||
let account_id = authenticate(&mut client, &config.auth)?;
|
||||
@@ -389,6 +390,7 @@ pub fn run(common: CommonConfig, config: ImapImportConfig) -> Result<Summary, Er
|
||||
compress: config.compress,
|
||||
policy,
|
||||
backoff: backoff.clone(),
|
||||
logger,
|
||||
},
|
||||
config.imap_connections.max(1),
|
||||
)
|
||||
|
||||
@@ -17,6 +17,7 @@ use crate::imap::name::encode_mailbox_name_with;
|
||||
use crate::imap::response::Untagged;
|
||||
use crate::imap::retry::{BackoffState, Disposition, RetryPolicy, classify};
|
||||
use crate::imap::transport::Connector;
|
||||
use crate::logging::Logger;
|
||||
|
||||
use super::coordinator::{Endpoint, ImapAuth, authenticate_client};
|
||||
use super::fetch::FetchAttrs;
|
||||
@@ -51,6 +52,7 @@ pub struct WorkerArgs {
|
||||
pub compress: bool,
|
||||
pub policy: RetryPolicy,
|
||||
pub backoff: BackoffState,
|
||||
pub logger: Logger,
|
||||
}
|
||||
|
||||
pub struct WorkerPool {
|
||||
@@ -208,6 +210,7 @@ fn connect_and_auth(args: &WorkerArgs) -> Result<ImapClient, ImapError> {
|
||||
&args.endpoint.host,
|
||||
args.endpoint.port,
|
||||
args.mode,
|
||||
args.logger,
|
||||
)?;
|
||||
authenticate_client(&mut client, &args.auth)
|
||||
.map_err(|e| ImapError::AuthFailed(e.to_string()))?;
|
||||
|
||||
@@ -65,6 +65,7 @@ fn reconnect(client: &mut SieveClient, ctx: &ControlCtx) -> Result<(), SieveErro
|
||||
ctx.port,
|
||||
ctx.mode,
|
||||
ctx.allow_cleartext,
|
||||
ctx.logger,
|
||||
)?;
|
||||
*client = new;
|
||||
match do_authenticate(client, &ctx.auth) {
|
||||
@@ -128,6 +129,7 @@ pub fn run(common: CommonConfig, config: ManageSieveImportConfig) -> Result<Summ
|
||||
endpoint.port,
|
||||
mode,
|
||||
config.allow_cleartext,
|
||||
logger,
|
||||
)
|
||||
.map_err(|e| Error::Connection(e.to_string()))?;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user