From 8ea1f46b109479185eb8c8577290e2f98d76cfd7 Mon Sep 17 00:00:00 2001 From: fylorn <249551762+fylorn@users.noreply.github.com> Date: Mon, 28 Sep 2026 17:55:36 +0800 Subject: [PATCH] test: hold the early-cancel client once its key is known, not on a timer a_client_that_leaves_before_routing_ends_leaves_at_most_one_cancelled_row sent five streamed requests on a 5 ms client timeout and needed at least one of them to be recorded. The guard is armed only once the API-key middleware has looked the key up; a client that leaves before then is nobody yet and, by design, writes nothing. When all five left that early there was no row, and the test failed with "5 left: []". Reproduced every time with a 1 ms timeout. The timer raced the other way too: at 20 ms, four of five clients had their stream started before leaving, rows the test did not count as clients that left. The test now holds the request where the key is known and no route is picked yet: it locks rbac_role_assignments, which the middleware reads right after arming the guard, waits until a query is blocked on that lock, and only then drops the client. It expects exactly one cancelled row (499, client_cancelled, cancelled_before: response) while the request is still held, and still only that row once the lock is let go. Co-Authored-By: Claude Opus 5.5 --- crates/test-support/tests/early_cancel.rs | 94 ++++++++++++++--------- 1 file changed, 57 insertions(+), 37 deletions(-) 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