postjevsql.git / tests / retries.rs

Retries follow the vendor SDKs (contract Execution): 408, 429 and 5xx are retried at most twice, a server delay is honoured, and each retry says which one it is. Any other status fails at once.

5use std::sync::atomic::{AtomicUsize, Ordering};
6use std::time::{Duration, Instant};
8use serde_json::json;
9use tokio_postgres::error::SqlState;
10use support::mock_jev::{MockJev, Reply};
11use support::{jev_instance, jev_instance_with, noul};
12
13const SELECT: &str = "SELECT jev_prob(t, 'Urgent?') FROM tickets t";
14
15#[tokio::test(flavor = "multi_thread")]
16async fn a_429_is_retried_after_the_servers_delay() {
17    let calls = AtomicUsize::new(0);
18    let mock = MockJev::start(move |_| match calls.fetch_add(1, Ordering::SeqCst) {
19        0 => Reply::json(429, json!({ "detail": { "error_type": "rate_limited", "message": "slow down" } }))
20            .header("retry-after-ms", "300"),
21        _ => Reply::json(200, noul(0.95)),
22    })
23    .await;
24    let (_pg, client) = jev_instance(&mock).await;
25
26    let start = Instant::now();
27    let row = client.query_one(SELECT, &[]).await.expect("retried to success");
28    assert_eq!(row.get::<_, f64>(0), 0.95);
29    assert!(start.elapsed() >= Duration::from_millis(300), "waited {:?}", start.elapsed());
30
31    let requests = mock.requests();
32    assert_eq!(requests.len(), 2);
33    assert!(!requests[0].headers.contains_key("x-typesafe-retry-count"));
34    assert_eq!(requests[1].headers["x-typesafe-retry-count"], "1");
35    assert_eq!(requests[0].body, requests[1].body, "a retry resends the same bytes");
36}
37
38#[tokio::test(flavor = "multi_thread")]
39async fn a_5xx_is_retried_twice_then_fails() {
40    let mock = MockJev::start(|_| {
41        Reply::json(503, json!({ "detail": { "error_type": "unavailable", "message": "down" } }))
42            .header("x-typesafe-request-id", "req-503")
43    })
44    .await;
45    let (_pg, client) = jev_instance(&mock).await;
46
47    let start = Instant::now();
48    let err = client.query_one(SELECT, &[]).await.expect_err("gives up");
49    let db = err.as_db_error().expect("an ERROR");
50    assert_eq!(mock.requests().len(), 3, "first attempt plus two retries");
51    // Backoff 0.5 s then 1 s, each less up to 25% jitter.
52    assert!(start.elapsed() >= Duration::from_millis(1125), "waited {:?}", start.elapsed());
53    assert_eq!(db.code(), &SqlState::EXTERNAL_ROUTINE_EXCEPTION);
54    assert!(db.message().contains("503") && db.message().contains("req-503"), "{}", db.message());
55    assert_eq!(db.detail(), Some("request id req-503; 3 attempts"));
56}
57
58#[tokio::test(flavor = "multi_thread")]
59async fn other_4xx_fail_at_once() {
60    let mock = MockJev::start(|_| {
61        Reply::json(403, json!({ "detail": { "error_type": "forbidden", "message": "no key" } }))
62            .header("x-typesafe-request-id", "req-403")
63    })
64    .await;
65    let (_pg, client) = jev_instance(&mock).await;
66
67    let err = client.query_one(SELECT, &[]).await.expect_err("fails");
68    let db = err.as_db_error().expect("an ERROR");
69    assert_eq!(mock.requests().len(), 1);
70    assert_eq!(db.code(), &SqlState::INVALID_AUTHORIZATION_SPECIFICATION);
71    assert!(db.message().contains("403") && db.message().contains("no key"), "{}", db.message());
72    assert_eq!(db.detail(), Some("request id req-403; 1 attempt"));
73}
74
75#[tokio::test(flavor = "multi_thread")]
76async fn a_connection_the_server_closed_is_redialled_not_retried() {
77    let mock = MockJev::start(|_| Reply::json(200, noul(0.95))).await;
78    let (_pg, client) = jev_instance(&mock).await;
79    client.query_one(SELECT, &[]).await.expect("first");
80
81    // The server closes the idle connection between statements.
82    mock.goaway();
83    tokio::time::sleep(Duration::from_millis(200)).await;
84
85    // A new namespace, so the cache cannot answer it.
86    client.batch_execute("SET jev.cache_namespace = 'second'").await.unwrap();
87    client.query_one(SELECT, &[]).await.expect("second, on a new connection");
88    let requests = mock.requests();
89    assert_eq!(requests.len(), 2, "the dead connection was never written to");
90    assert!(!requests[1].headers.contains_key("x-typesafe-retry-count"), "not a retry");
91    assert_eq!(mock.connections(), 2);
92}
93
94#[tokio::test(flavor = "multi_thread")]
95async fn an_untrusted_certificate_is_not_retried() {
96    let mock = MockJev::start(|_| Reply::json(200, noul(0.95))).await;
97    // Without jev.ca_file the mock's throwaway CA is not trusted.
98    let (_pg, client) = jev_instance_with(&mock, &[("jev.ca_file", "")]).await;
99
100    let err = client.query_one(SELECT, &[]).await.expect_err("refused");
101    let db = err.as_db_error().unwrap();
102    assert_eq!(db.code(), &SqlState::SQLCLIENT_UNABLE_TO_ESTABLISH_SQLCONNECTION);
103    assert!(db.message().contains("TLS"), "{}", db.message());
104    assert_eq!(mock.accepts(), 1, "dialled once: not retried, not redialled");
105    assert_eq!(mock.connections(), 0, "the handshake never completed");
106}
107
108#[tokio::test(flavor = "multi_thread")]
109async fn an_oversized_answer_is_not_asked_for_again() {
110    let big = "x".repeat(9 << 20);
111    let mock = MockJev::start(move |_| Reply::json(200, json!({ "padding": big }))).await;
112    let (_pg, client) = jev_instance(&mock).await;
113
114    let err = client.query_one(SELECT, &[]).await.expect_err("too large");
115    let db = err.as_db_error().unwrap();
116    assert_eq!(db.code(), &SqlState::PROGRAM_LIMIT_EXCEEDED);
117    assert_eq!(mock.requests().len(), 1, "already billed once; never again");
118}