diff --git a/crates/test-support/tests/early_cancel.rs b/crates/test-support/tests/early_cancel.rs index 81a715be..ecac41ed 100644 --- a/crates/test-support/tests/early_cancel.rs +++ b/crates/test-support/tests/early_cancel.rs @@ -112,54 +112,74 @@ async fn a_client_that_leaves_while_a_whole_answer_is_awaited_is_recorded() { assert_eq!(detail["model_id"], "early-cancel-whole", "{detail}"); } -/// Leaving at once: the request is dropped somewhere in the key's -/// roles, the limits or routing. A client that left before its key was -/// even looked up is nobody yet and writes nothing; every other one -/// leaves exactly one cancelled row. Never a success, never two. +/// Leaving before routing ends: the middleware loads the key's roles +/// right after it arms the guard, so a request held there is known but +/// not yet routed. A lock on the role assignments holds it, and the +/// client goes while it waits. One cancelled row, written as it leaves, +/// and nothing more once the lock is let go. Never a success, never two. +/// +/// A client that leaves before its key is even looked up is nobody yet +/// and writes nothing. Five clients on a 5 ms timer used to stand in for +/// this one, and on a slow runner all five could leave that early. #[ignore = "integration test — run via `make test-it`"] #[tokio::test] async fn a_client_that_leaves_before_routing_ends_leaves_at_most_one_cancelled_row() { let app = TestApp::spawn_with_clickhouse().await; - let (server, user_id, key) = slow_route(&app, "early-cancel-quick").await; + let (_server, user_id, key) = slow_route(&app, "early-cancel-quick").await; - let client = reqwest::Client::builder() - .timeout(std::time::Duration::from_millis(5)) - .build() + // Nothing else in this test reads role assignments, so whatever waits + // on the lock is the request. + let mut roles = app.db.begin().await.unwrap(); + sqlx::query("LOCK TABLE rbac_role_assignments IN ACCESS EXCLUSIVE MODE") + .execute(&mut *roles) + .await + .unwrap(); + let holder: i32 = sqlx::query_scalar("SELECT pg_backend_pid()") + .fetch_one(&mut *roles) + .await .unwrap(); - let mut left = 0; - for _ in 0..5 { - let r = client - .post(format!("{}/v1/chat/completions", app.gateway_url)) - .bearer_auth(&key) + + let url = format!("{}/v1/chat/completions", app.gateway_url); + let call = tokio::spawn(async move { + reqwest::Client::new() + .post(url) + .bearer_auth(key) .json(&json!({"model": "early-cancel-quick", "stream": true, "messages": [{"role": "user", "content": "hi"}]})) .send() - .await; - if r.is_err() { - left += 1; - } - } - assert!(left > 0, "a 5 ms client never left early"); - - // Let the audit pipeline flush whatever was written. - let found = rows(&app, user_id, left).await; - // At most one row per client that left; and some of them left after - // the key was known — before, none of these were ever recorded. - assert!( - !found.is_empty() && found.len() <= left, - "{left} left: {found:?}" - ); - for row in &found { - let (status, ..) = row; - if *status == 499 { - // A stream that had started carries no `cancelled_before`. - let detail: Value = serde_json::from_str(&row.3).unwrap(); - assert_eq!(detail["stream_outcome"], "client_cancelled", "{detail}"); - } else { - panic!("a request whose client left was logged as {status}: {row:?}"); + .await + }); + let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(30); + loop { + let waiting: bool = sqlx::query_scalar( + "SELECT EXISTS (SELECT 1 FROM pg_stat_activity \ + WHERE $1 = ANY(pg_blocking_pids(pid)))", + ) + .bind(holder) + .fetch_one(&app.db) + .await + .unwrap(); + if waiting { + break; } + assert!( + !call.is_finished() && tokio::time::Instant::now() < deadline, + "the request never reached the key's roles" + ); + tokio::time::sleep(std::time::Duration::from_millis(10)).await; } - drop(server); + call.abort(); + + // Written as the client left: the request is still held. + let found = rows(&app, user_id, 1).await; + assert_eq!(found.len(), 1, "{found:?}"); + assert_cancelled(&found[0]); + + // Once the lock is let go, a request that outlived its client would + // carry on and record again. + roles.rollback().await.unwrap(); + tokio::time::sleep(std::time::Duration::from_secs(3)).await; + assert_eq!(rows(&app, user_id, 1).await, found); } /// A stream that started records its own cancel; the guard is disarmed