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).
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.
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.
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}