export: batch Email/import, upload in parallel, and show progress #10
@@ -97,6 +97,10 @@ impl Inner {
|
|||||||
#[derive(Debug, Clone, Copy)]
|
#[derive(Debug, Clone, Copy)]
|
||||||
enum Kind {
|
enum Kind {
|
||||||
Api,
|
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,
|
Upload,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -220,6 +224,23 @@ impl HttpClient {
|
|||||||
.map_err(|e| JmapError::Malformed(format!("response is not valid json: {e}")))
|
.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<Value, JmapError> {
|
||||||
|
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(
|
pub fn upload(
|
||||||
&self,
|
&self,
|
||||||
upload_url: &str,
|
upload_url: &str,
|
||||||
@@ -298,6 +319,16 @@ impl HttpClient {
|
|||||||
body: truncate(&body),
|
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 => {
|
StatusOutcome::Retryable => {
|
||||||
attempt += 1;
|
attempt += 1;
|
||||||
self.inner.retries_total.fetch_add(1, Ordering::Relaxed);
|
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) {
|
Attempt::Transport(err) => match transport_disposition(&err) {
|
||||||
Disposition::Fatal => return Err(err),
|
Disposition::Fatal => return Err(err),
|
||||||
Disposition::Retryable => {
|
Disposition::Retryable => {
|
||||||
@@ -659,6 +691,90 @@ mod tests {
|
|||||||
use super::*;
|
use super::*;
|
||||||
use crate::net::CertOverride;
|
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<AtomicU64>) {
|
||||||
|
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]
|
#[test]
|
||||||
fn invalid_certificates_are_accepted_only_for_the_named_host() {
|
fn invalid_certificates_are_accepted_only_for_the_named_host() {
|
||||||
let client = HttpClient::new(
|
let client = HttpClient::new(
|
||||||
|
|||||||
@@ -131,6 +131,14 @@ impl Request {
|
|||||||
let value = client.post_json(api_url, &self.envelope()?)?;
|
let value = client.post_json(api_url, &self.envelope()?)?;
|
||||||
Response::parse(value)
|
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<Response, JmapError> {
|
||||||
|
let value = client.post_json_once(api_url, &self.envelope()?)?;
|
||||||
|
Response::parse(value)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
|
|||||||
Reference in New Issue
Block a user