diff --git a/CHANGELOG.md b/CHANGELOG.md index 65ae5e4..ee2a1be 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,7 @@ All notable changes to this project will be documented in this file. This projec ### Changed ### Fixed +- Improve verbosity (#4 #19). - WebDAV import materialised the account root collection as a directory named after the account displayname (#18). - Report user friendly error message when `urn:ietf:params:jmap:principals` is not supported and no accountId is provided (#21). - Report which email failed to import when the blob is too large (#22). diff --git a/src/cli.rs b/src/cli.rs index ccccf87..457da5c 100644 --- a/src/cli.rs +++ b/src/cli.rs @@ -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, diff --git a/src/dav/client.rs b/src/dav/client.rs index e6138ce..2766060 100644 --- a/src/dav/client.rs +++ b/src/dav/client.rs @@ -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) + } } } diff --git a/src/exchange_ews/client.rs b/src/exchange_ews/client.rs index cebb9c9..6e69fc8 100644 --- a/src/exchange_ews/client.rs +++ b/src/exchange_ews/client.rs @@ -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) + } } } } diff --git a/src/exchange_graph/client.rs b/src/exchange_graph/client.rs index 10ebd3c..104f61c 100644 --- a/src/exchange_graph/client.rs +++ b/src/exchange_graph/client.rs @@ -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}")), diff --git a/src/imap/client.rs b/src/imap/client.rs index 89ae2a6..eb374cd 100644 --- a/src/imap/client.rs +++ b/src/imap/client.rs @@ -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>, @@ -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 { let stream: Box = 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, } +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 { 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), } } diff --git a/src/jmap/http.rs b/src/jmap/http.rs index 19f2bc7..eaa0db3 100644 --- a/src/jmap/http.rs +++ b/src/jmap/http.rs @@ -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]; diff --git a/src/logging.rs b/src/logging.rs index 5903f5c..e17961d 100644 --- a/src/logging.rs +++ b/src/logging.rs @@ -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 "".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::(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} "); + } + if first.eq_ignore_ascii_case("AUTHENTICATE") { + let mechanism = parts.next().unwrap_or(""); + return match parts.next() { + Some(_) => format!("{first} {mechanism} "), + 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 "); + } + + #[test] + fn redact_wire_hides_authenticate_initial_response() { + assert_eq!( + redact_wire("AUTHENTICATE PLAIN AGFsaWNlAHMzY3JldA=="), + "AUTHENTICATE PLAIN " + ); + } + + #[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), ""); + } + + #[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), + ); + } } diff --git a/src/managesieve/client.rs b/src/managesieve/client.rs index 619487c..71b32a1 100644 --- a/src/managesieve/client.rs +++ b/src/managesieve/client.rs @@ -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 { let stream: Box = 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 { + 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 { + 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 String { + let word = match status.status { + Status::Ok => "OK", + Status::No => "NO", + Status::Bye => "BYE", + }; + format!("{word} {}", status.text) +} + fn noop_stream() -> Box { Box::new(NoopStream) } diff --git a/src/sync/import_imap/coordinator.rs b/src/sync/import_imap/coordinator.rs index f44e077..062bc50 100644 --- a/src/sync/import_imap/coordinator.rs +++ b/src/sync/import_imap/coordinator.rs @@ -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 Result Result { &args.endpoint.host, args.endpoint.port, args.mode, + args.logger, )?; authenticate_client(&mut client, &args.auth) .map_err(|e| ImapError::AuthFailed(e.to_string()))?; diff --git a/src/sync/import_managesieve/coordinator.rs b/src/sync/import_managesieve/coordinator.rs index ac11cb3..71cee22 100644 --- a/src/sync/import_managesieve/coordinator.rs +++ b/src/sync/import_managesieve/coordinator.rs @@ -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 Result { "127.0.0.1", server.port, ConnectMode::Plain, + Logger::from_flags(false, 0), ) }