postjevsql.git / tests / cancel.rs

Giving up on a request resets its HTTP/2 stream, whether the attempt timed out or the query was cancelled (contract Testing: the reset is proven at the server, not inferred from the query returning).

5use std::sync::atomic::{AtomicUsize, Ordering};
6use std::time::{Duration, Instant};
8use support::mock_jev::{MockJev, Reply};
9use support::{jev_instance, noul};
10use tokio_postgres::error::SqlState;
11
12const SELECT: &str = "SELECT jev_prob(t, 'Urgent?') FROM tickets t";
13
14async fn until(what: &str, mut done: impl FnMut() -> bool) {
15    // A bound on a hang, not a measurement: generous for a loaded machine.
16    let deadline = Instant::now() + Duration::from_secs(15);
17    while !done() {
18        assert!(Instant::now() < deadline, "timed out waiting for {what}");
19        tokio::time::sleep(Duration::from_millis(20)).await;
20    }
21}
22
23#[tokio::test(flavor = "multi_thread")]
24async fn a_stalled_attempt_is_abandoned_after_10_s_and_retried() {
25    let calls = AtomicUsize::new(0);
26    let mock = MockJev::start(move |_| match calls.fetch_add(1, Ordering::SeqCst) {
27        0 => Reply::json(200, noul(0.5)).after(Duration::from_secs(60)),
28        _ => Reply::json(200, noul(0.95)),
29    })
30    .await;
31    let (_pg, client) = jev_instance(&mock).await;
32
33    let start = Instant::now();
34    let row = client.query_one(SELECT, &[]).await.expect("the retry answers");
35    assert_eq!(row.get::<_, f64>(0), 0.95);
36    assert!(start.elapsed() >= Duration::from_secs(10), "{:?}", start.elapsed());
37    assert_eq!(mock.requests().len(), 2);
38    until("the stalled stream's reset", || mock.abandoned() == 1).await;
39    let in_flight: i64 = client.query_one("SELECT in_flight FROM jev_stats()", &[]).await.unwrap().get(0);
40    assert_eq!(in_flight, 0, "the abandoned attempt left the count");
41}

How a test cancels: from a second session (with the first's pid), or with the first session's own cancel key.

45enum Cancel {
46    PgCancelBackend,
47    CancelRequest,
48}
50async fn cancelled_mid_request(how: Cancel) {
51    let mock = MockJev::start(|_| Reply::json(200, noul(0.5)).after(Duration::from_secs(60))).await;
52    let (pg, client) = jev_instance(&mock).await;
53    let pid: i32 = client.query_one("SELECT pg_backend_pid()", &[]).await.unwrap().get(0);
54    let token = client.cancel_token();
55
56    // The session outlives the cancel: a closed backend would close the
57    // socket, and the server would see that instead of a stream reset.
58    let query = tokio::spawn(async move {
59        let result = client.query_one(SELECT, &[]).await.map(drop);
60        (client, result)
61    });
62    until("the request to arrive", || mock.requests().len() == 1).await;
63
64    match how {
65        Cancel::PgCancelBackend => {
66            let (other, conn) = tokio_postgres::connect(&pg.conn_str(), tokio_postgres::NoTls).await.unwrap();
67            tokio::spawn(conn);
68            let ok: bool = other.query_one("SELECT pg_cancel_backend($1)", &[&pid]).await.unwrap().get(0);
69            assert!(ok);
70        }
71        Cancel::CancelRequest => token.cancel_query(tokio_postgres::NoTls).await.unwrap(),
72    }
73
74    let (client, result) = query.await.unwrap();
75    let err = result.expect_err("cancelled");
76    let db = err.as_db_error().expect("an ERROR");
77    assert_eq!(db.code(), &SqlState::QUERY_CANCELED);
78    // statement_timeout shares the SQLSTATE; the message tells them apart.
79    assert_eq!(db.message(), "canceling statement due to user request");
80    until("the stream's reset", || mock.abandoned() == 1).await;
81    assert_eq!(mock.requests().len(), 1, "a cancel is not retried");
82    assert_eq!(mock.connections(), 1);
83
84    // The connection survives: only the stream was reset.
85    client.query_one("SELECT 1", &[]).await.expect("the session is fine");
86    let in_flight: i64 = client.query_one("SELECT in_flight FROM jev_stats()", &[]).await.unwrap().get(0);
87    assert_eq!(in_flight, 0, "the cancelled stream left the count");
88}
89
90#[tokio::test(flavor = "multi_thread")]
91async fn pg_cancel_backend_resets_the_stream() {
92    cancelled_mid_request(Cancel::PgCancelBackend).await;
93}
94
95#[tokio::test(flavor = "multi_thread")]
96async fn a_protocol_cancel_resets_the_stream() {
97    cancelled_mid_request(Cancel::CancelRequest).await;
98}

Every fd the backend holds, by what it points at, from /proc (the backend runs as this user). All of them, not only sockets: the reactor's per-wait epoll set is an fd too.

103fn fds(pid: i32) -> Vec<String> {
104    let mut targets: Vec<String> = std::fs::read_dir(format!("/proc/{pid}/fd"))
105        .unwrap()
106        .filter_map(|e| std::fs::read_link(e.unwrap().path()).ok())
107        .map(|target| target.to_string_lossy().into_owned())
108        .collect();
109    targets.sort();
110    targets
111}

Sockets and anonymous inodes (epoll, eventfd) are distinct per open, so their inode numbers are stripped to compare kinds and counts.

115fn kinds(fds: &[String]) -> Vec<String> {
116    fds.iter()
117        .map(|t| match t.split_once(":[") {
118            Some((kind, _)) if kind == "socket" || kind == "pipe" => format!("{kind}:"),
119            _ => t.clone(),
120        })
121        .collect()
122}
124#[tokio::test(flavor = "multi_thread")]
125async fn a_hundred_cancels_leak_nothing() {
126    let mock = MockJev::start(|_| Reply::json(200, noul(0.5)).after(Duration::from_secs(60))).await;
127    let (pg, client) = jev_instance(&mock).await;
128    let client = std::sync::Arc::new(client);
129    let pid: i32 = client.query_one("SELECT pg_backend_pid()", &[]).await.unwrap().get(0);
130    let token = client.cancel_token();
131    let (other, conn) = tokio_postgres::connect(&pg.conn_str(), tokio_postgres::NoTls).await.unwrap();
132    tokio::spawn(conn);
133
134    // Both ways to cancel, alternating: a query cancelled by one must
135    // leave nothing the other would then trip over.
136    let mut baseline = None;
137    for round in 1..=100 {
138        let c = client.clone();
139        let query = tokio::spawn(async move { c.query_one(SELECT, &[]).await.map(drop) });
140        until("the request to arrive", || mock.requests().len() == round).await;
141        if round % 2 == 0 {
142            let ok: bool = other.query_one("SELECT pg_cancel_backend($1)", &[&pid]).await.unwrap().get(0);
143            assert!(ok, "round {round}");
144        } else {
145            token.cancel_query(tokio_postgres::NoTls).await.unwrap();
146        }
147        let err = query.await.unwrap().expect_err("cancelled");
148        assert_eq!(err.as_db_error().unwrap().code(), &SqlState::QUERY_CANCELED, "round {round}");
149        // After the first round: the connection and the relation's files
150        // are open by then, and are meant to stay open.
151        if round == 1 {
152            baseline = Some(fds(pid));
153        }
154    }
155    until("every reset", || mock.abandoned() == 100).await;
156    assert_eq!(mock.requests().len(), 100, "a cancel is not retried");
157    assert_eq!(mock.connections(), 1, "one connection for all hundred");
158    let (before, after) = (baseline.unwrap(), fds(pid));
159    // The instrument sees the client socket and the h2 connection.
160    assert!(before.iter().filter(|t| t.starts_with("socket:")).count() >= 2, "{before:#?}");
161    assert_eq!(kinds(&after), kinds(&before), "no fds leaked: {before:#?} then {after:#?}");
162    let in_flight: i64 = client.query_one("SELECT in_flight FROM jev_stats()", &[]).await.unwrap().get(0);
163    assert_eq!(in_flight, 0, "every cancelled stream left the count");
164}
165
166#[tokio::test(flavor = "multi_thread")]
167async fn terminating_a_backend_mid_request_crashes_nothing() {
168    let mock = MockJev::start(|_| Reply::json(200, noul(0.5)).after(Duration::from_secs(60))).await;
169    let (pg, client) = jev_instance(&mock).await;
170    let pid: i32 = client.query_one("SELECT pg_backend_pid()", &[]).await.unwrap().get(0);
171    let query = tokio::spawn(async move {
172        let result = client.query_one(SELECT, &[]).await.map(drop);
173        (client, result)
174    });
175    until("the request to arrive", || mock.requests().len() == 1).await;
176
177    // A second session that must survive: a backend that aborted (a
178    // panic in a destructor during FATAL exit) makes the postmaster
179    // restart every session.
180    let (other, conn) = tokio_postgres::connect(&pg.conn_str(), tokio_postgres::NoTls).await.unwrap();
181    tokio::spawn(conn);
182    let ok: bool = other.query_one("SELECT pg_terminate_backend($1)", &[&pid]).await.unwrap().get(0);
183    assert!(ok);
184    let (_client, result) = query.await.unwrap();
185    assert!(result.is_err(), "the terminated session's query fails");
186    until("the stream's reset", || mock.abandoned() == 1).await;
187
188    tokio::time::sleep(Duration::from_millis(300)).await;
189    let alive: i32 = other.query_one("SELECT 1", &[]).await.expect("no crash-restart").get(0);
190    assert_eq!(alive, 1);
191}