diff --git a/src/jmap/http.rs b/src/jmap/http.rs index 4623c46..3c9e614 100644 --- a/src/jmap/http.rs +++ b/src/jmap/http.rs @@ -97,6 +97,10 @@ impl Inner { #[derive(Debug, Clone, Copy)] enum Kind { Api, + /// An API call that must not be sent twice: a write the server may + /// already have applied when the connection failed. A transport error is + /// returned to the caller, which checks the target instead of resending. + ApiOnce, Upload, } @@ -220,6 +224,23 @@ impl HttpClient { .map_err(|e| JmapError::Malformed(format!("response is not valid json: {e}"))) } + /// As `post_json`, but a transport failure is not retried: the request + /// may have reached the server, and sending it again could apply it twice. + /// Throttling and `503` answers, which mean it was not processed, are still + /// retried. + pub fn post_json_once(&self, url: &str, body: &Value) -> Result { + let payload = serde_json::to_vec(body)?; + let raw = self.execute( + Kind::ApiOnce, + "POST", + url, + Some(&payload), + Some("application/json"), + )?; + serde_json::from_slice(&raw) + .map_err(|e| JmapError::Malformed(format!("response is not valid json: {e}"))) + } + pub fn upload( &self, upload_url: &str, @@ -298,6 +319,16 @@ impl HttpClient { body: truncate(&body), }); } + // A gateway error may come back after the server + // behind it applied the write: not safe to resend. + StatusOutcome::Retryable + if matches!(kind, Kind::ApiOnce) && matches!(status, 502 | 504) => + { + return Err(JmapError::HttpStatus { + status, + body: truncate(&body), + }); + } StatusOutcome::Retryable => { attempt += 1; self.inner.retries_total.fetch_add(1, Ordering::Relaxed); @@ -344,6 +375,7 @@ impl HttpClient { } } } + Attempt::Transport(err) if matches!(kind, Kind::ApiOnce) => return Err(err), Attempt::Transport(err) => match transport_disposition(&err) { Disposition::Fatal => return Err(err), Disposition::Retryable => { @@ -659,6 +691,90 @@ mod tests { use super::*; use crate::net::CertOverride; + /// A server that reads each request and hangs up without answering: a + /// transport failure after the request has been sent. Returns its URL and + /// the count of requests it has seen. + fn hang_up_server() -> (String, Arc) { + use std::io::Read; + let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let url = format!("http://{}/jmap/api", listener.local_addr().unwrap()); + let seen = Arc::new(AtomicU64::new(0)); + let counter = seen.clone(); + std::thread::spawn(move || { + for stream in listener.incoming() { + let Ok(mut stream) = stream else { break }; + let mut buf = [0u8; 4096]; + let _ = stream.read(&mut buf); + counter.fetch_add(1, Ordering::SeqCst); + } + }); + (url, seen) + } + + fn quick_retries(max_retries: u32) -> RetryPolicy { + RetryPolicy { + max_retries, + base: Duration::from_millis(1), + cap: Duration::from_millis(2), + } + } + + #[test] + fn a_write_sent_once_is_not_resent_after_a_transport_failure() { + let (url, seen) = hang_up_server(); + let client = HttpClient::new( + Auth::Bearer { token: "t".into() }, + quick_retries(2), + CertOverride::none(), + ); + let err = client + .post_json_once(&url, &serde_json::json!({})) + .unwrap_err(); + assert!(matches!(err, JmapError::Transport(_)), "{err}"); + assert_eq!(seen.load(Ordering::SeqCst), 1, "sent exactly once"); + + let (url, seen) = hang_up_server(); + let _ = client.post_json(&url, &serde_json::json!({})); + assert_eq!( + seen.load(Ordering::SeqCst), + 3, + "an ordinary call retries twice" + ); + } + + #[test] + fn a_write_sent_once_is_not_resent_after_a_gateway_timeout_but_is_after_503() { + let mut server = mockito::Server::new(); + let url = format!("{}/jmap/api", server.url()); + let client = HttpClient::new( + Auth::Bearer { token: "t".into() }, + quick_retries(2), + CertOverride::none(), + ); + let gateway = server + .mock("POST", "/jmap/api") + .with_status(504) + .expect(1) + .create(); + let err = client + .post_json_once(&url, &serde_json::json!({})) + .unwrap_err(); + assert!( + matches!(err, JmapError::HttpStatus { status: 504, .. }), + "{err}" + ); + gateway.assert(); + gateway.remove(); + + let busy = server + .mock("POST", "/jmap/api") + .with_status(503) + .expect(3) + .create(); + let _ = client.post_json_once(&url, &serde_json::json!({})); + busy.assert(); + } + #[test] fn invalid_certificates_are_accepted_only_for_the_named_host() { let client = HttpClient::new( diff --git a/src/jmap/request.rs b/src/jmap/request.rs index 7a6b4cc..9c1265f 100644 --- a/src/jmap/request.rs +++ b/src/jmap/request.rs @@ -131,6 +131,14 @@ impl Request { let value = client.post_json(api_url, &self.envelope()?)?; Response::parse(value) } + + /// As `send`, for a write that must not be applied twice: a transport + /// failure comes back as an error instead of being resent (see + /// `HttpClient::post_json_once`). + pub fn send_once(&self, client: &HttpClient, api_url: &str) -> Result { + let value = client.post_json_once(api_url, &self.envelope()?)?; + Response::parse(value) + } } #[derive(Debug)]