From 6e43f4dca8b20862a95c9adb72c4008c6b1ec18d Mon Sep 17 00:00:00 2001 From: Jared Lunde Date: Sun, 4 Oct 2026 12:46:38 -0700 Subject: [PATCH 1/2] gateway tests: make timing assertions hold on a heavily loaded host Wall-clock assertions failed at load averages of 60-230 without any code getting slower. Each fragile one now asserts something load cannot fake or break: - D217 merge test: thread CPU time (rustix clock_gettime, dev-dep; the crate forbids unsafe) and scaling, 4x the messages < 8x the CPU, in place of < 2 s wall. The reintroduced quadratic merge reads 15.7x. - deadline, breaker window, cache TTL, key cooldown unit tests: measured elapsed time, a stepped clock, or re-running a run that a stall made inconclusive, in place of fixed sleeps with an implicit upper bound. - streaming claims: the scripted provider writes the next event only once the client holds the last one (new Step::Until), so the ordering is the proof rather than inter-arrival gaps. - request deadlines, header stall, dead h2 PING, early answer, large-body failover: the bound is pushed far below the alternative ending, which is pushed far out (a trickled body that lasts over a minute), and the cause is asserted from the answer, row, metric or drain log line. - SEC-17 deny bound: 2 s on an idle host, stretched by 40 round trips measured in the same run when the host is loaded (common::stretched). - rate-limit and log-cap floods: bursts re-sent until one provably landed inside the windows the assertion needs. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01JimHGjsfk2Ktm5GxyZJKKk --- Cargo.lock | 1 + crates/gateway/Cargo.toml | 4 + crates/gateway/src/cache.rs | 28 +++- crates/gateway/src/circuit_breaker.rs | 16 +- crates/gateway/src/deadline.rs | 43 ++++- crates/gateway/src/route.rs | 15 +- crates/gateway/src/translate_request_tests.rs | 61 +++++-- crates/gateway/tests/claims_security.rs | 83 +++++++--- crates/gateway/tests/claims_streaming.rs | 149 ++++++++++-------- crates/gateway/tests/common/mod.rs | 42 +++++ crates/gateway/tests/large_bodies.rs | 11 +- crates/gateway/tests/reliability_breaker.rs | 5 +- .../tests/reliability_early_response.rs | 6 +- crates/gateway/tests/reliability_lifecycle.rs | 41 +++-- crates/gateway/tests/reliability_log_stall.rs | 39 ++++- crates/gateway/tests/replicas.rs | 28 +++- crates/gateway/tests/request_deadlines.rs | 53 +++++-- 17 files changed, 455 insertions(+), 170 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index ab7f9a68..4b6db53f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -388,6 +388,7 @@ dependencies = [ "reqwest", "ring", "rustc-hash", + "rustix", "rustls", "serde", "serde_json", diff --git a/crates/gateway/Cargo.toml b/crates/gateway/Cargo.toml index 1ee6b47e..d9e828ce 100644 --- a/crates/gateway/Cargo.toml +++ b/crates/gateway/Cargo.toml @@ -89,6 +89,10 @@ hyper-util = { version = "0.1", features = ["tokio", "server-auto"] } # fork/timeout runner, so no rusty-fork or tempfile in the build graph. proptest = { version = "1.10", default-features = false, features = ["std"] } rcgen = "0.13" +# Thread CPU time (`clock_gettime(CLOCK_THREAD_CPUTIME_ID)`) behind a safe API, for perf assertions +# that must hold on a loaded host (`unsafe_code` is forbidden, so not raw `libc`). Already in the +# build graph via `tempfile` and `clap`: this only adds its `time` feature. +rustix = { version = "1", default-features = false, features = ["std", "time"] } reqwest = { version = "0.13", default-features = false, features = ["json", "rustls", "http2"] } tokio-rustls = "0.26" diff --git a/crates/gateway/src/cache.rs b/crates/gateway/src/cache.rs index 891c5976..e7a15a22 100644 --- a/crates/gateway/src/cache.rs +++ b/crates/gateway/src/cache.rs @@ -484,17 +484,29 @@ mod tests { #[test] fn insert_drops_an_expired_prefix_without_evicting_a_newer_live_entry() { - let c = ResponseCache::new(Duration::from_millis(40), 2, 1024); + const TTL: Duration = Duration::from_millis(40); let k1 = key_of(1, "/v1", b"a", &[]); let k2 = key_of(1, "/v1", b"b", &[]); - c.insert(k1, entry(b"old")); - std::thread::sleep(Duration::from_millis(50)); - c.insert(k2, entry(b"new")); let k3 = key_of(1, "/v1", b"c", &[]); - c.insert(k3, entry(b"newer")); - assert!(c.get(&k1).is_none()); - assert_eq!(c.get(&k2).unwrap().body.as_ref(), b"new"); - assert_eq!(c.get(&k3).unwrap().body.as_ref(), b"newer"); + // k2 and k3 must still be live when read back. A loaded host can stall this thread past + // the TTL between the inserts and the reads; such a run proves nothing, so it is re-run. + for _ in 0..50 { + let c = ResponseCache::new(TTL, 2, 1024); + c.insert(k1, entry(b"old")); + std::thread::sleep(TTL + Duration::from_millis(10)); + let live_from = Instant::now(); + c.insert(k2, entry(b"new")); + c.insert(k3, entry(b"newer")); + let (got1, got2, got3) = (c.get(&k1), c.get(&k2), c.get(&k3)); + if live_from.elapsed() >= TTL { + continue; + } + assert!(got1.is_none()); + assert_eq!(got2.unwrap().body.as_ref(), b"new"); + assert_eq!(got3.unwrap().body.as_ref(), b"newer"); + return; + } + panic!("every run stalled past the {TTL:?} TTL"); } #[test] diff --git a/crates/gateway/src/circuit_breaker.rs b/crates/gateway/src/circuit_breaker.rs index 10d79a74..84f18d48 100644 --- a/crates/gateway/src/circuit_breaker.rs +++ b/crates/gateway/src/circuit_breaker.rs @@ -917,15 +917,23 @@ mod tests { #[test] fn test_windowed_resets_after_window() { - // Note: window uses second-level precision, so use 1 second window - let cb = CircuitBreaker::new(CircuitBreakerConfig::windowed(3, Duration::from_secs(1))); + // A hand-stepped clock: on the real one, a second boundary falling between two of these + // calls (a preempted thread, a loaded host) would split one window into two. + static NOW: AtomicU64 = AtomicU64::new(100); + fn clock() -> u64 { + NOW.load(Ordering::Relaxed) + } + let cb = CircuitBreaker::with_clock( + CircuitBreakerConfig::windowed(3, Duration::from_secs(1)), + clock, + ); cb.record_failure(); cb.record_failure(); assert_eq!(cb.state(), CircuitState::Closed { failure_count: 2 }); - // Wait for window to expire (1 second + buffer) - thread::sleep(Duration::from_millis(1100)); + // The window (second-level precision) expires. + NOW.store(101, Ordering::Relaxed); // This failure starts a new window cb.record_failure(); diff --git a/crates/gateway/src/deadline.rs b/crates/gateway/src/deadline.rs index 8f58c7d4..f1c6452d 100644 --- a/crates/gateway/src/deadline.rs +++ b/crates/gateway/src/deadline.rs @@ -99,8 +99,14 @@ mod tests { let start = Instant::now(); let d = at(start, 2); let left = remaining(d).unwrap(); + // Measured after the read, so it covers however long this thread was descheduled: what is + // left plus what has passed is the whole two seconds, less the millisecond truncation. + let passed = start.elapsed(); assert!(left <= Duration::from_secs(2), "{left:?}"); - assert!(left > Duration::from_millis(1500), "{left:?}"); + assert!( + left + passed + Duration::from_millis(1) >= Duration::from_secs(2), + "{left:?} left after {passed:?}" + ); assert_eq!( remaining(0), Some(Duration::ZERO), @@ -108,12 +114,39 @@ mod tests { ); } + /// The coarse clock advances on its own and expires a deadline, never early: it lags real time + /// (by up to a tick, more on a loaded host), so a deadline it calls expired has truly passed. + /// Polled rather than slept on, so a ticker thread a loaded host runs late cannot fail it. #[test] fn the_coarse_clock_ticks_and_expires_deadlines() { start(); - let d = at(Instant::now(), 1); - assert!(!expired(d)); - std::thread::sleep(Duration::from_millis(2300)); - assert!(expired(d), "now {} deadline {d}", now_ms()); + let begun = Instant::now(); + let d = at(begun, 1); + let mut seen = now_ms(); + loop { + let now = now_ms(); + if now != seen { + // A tick. The ticker sleeps a whole TICK between stores, and load only adds to it. + assert!( + now + 1 >= seen + TICK.as_millis() as u64, + "ticked {seen} -> {now}" + ); + seen = now; + } + if now >= d { + break; + } + assert!( + begun.elapsed() < Duration::from_secs(30), + "the coarse clock never reached the deadline: now {now} deadline {d}" + ); + std::thread::sleep(Duration::from_millis(5)); + } + let passed = begun.elapsed(); + assert!(expired(d)); + assert!( + passed + Duration::from_millis(1) >= Duration::from_secs(1), + "expired after only {passed:?}" + ); } } diff --git a/crates/gateway/src/route.rs b/crates/gateway/src/route.rs index 0687ce28..78647a1c 100644 --- a/crates/gateway/src/route.rs +++ b/crates/gateway/src/route.rs @@ -1237,11 +1237,16 @@ mod tests { ProviderMetrics::disconnected(), None, ); - // A cooldown that ends 200 ms from now, rather than `KEY_COOLDOWN`'s minute. - p.pool_auth[0] - .bad_until_ms - .store(clock_ms().saturating_add(200), Ordering::Relaxed); - assert_eq!(p.first_key(), 1); + // A cooldown that ends 200 ms from now, rather than `KEY_COOLDOWN`'s minute. It only shows + // the key cooling if it is read before those 200 ms are up; a run a loaded host stalled + // past them proves nothing, so it is re-run. + let cooling_read = (0..50).find_map(|_| { + let until = clock_ms().saturating_add(200); + p.pool_auth[0].bad_until_ms.store(until, Ordering::Relaxed); + let first = p.first_key(); + (clock_ms() < until).then_some(first) + }); + assert_eq!(cooling_read, Some(1), "the cooling key is skipped"); std::thread::sleep(Duration::from_millis(250)); assert_eq!(p.first_key(), 0, "the cooldown has passed"); } diff --git a/crates/gateway/src/translate_request_tests.rs b/crates/gateway/src/translate_request_tests.rs index 8039c538..fcc38cc7 100644 --- a/crates/gateway/src/translate_request_tests.rs +++ b/crates/gateway/src/translate_request_tests.rs @@ -1696,24 +1696,30 @@ fn consecutive_user_messages_keep_their_text_apart() { /// client split across messages) merges in time linear in its length, and every message keeps its /// block. Each merge used to copy every block gathered so far: 256 KiB of user turns took 4.4 s of /// CPU in a release build, 1 MiB over a minute, all on one request. +/// +/// Asserted as scaling, in this thread's CPU time, so a loaded host cannot fail it: quadrupling the +/// run must cost under 8x (linear is ~4x; the quadratic merge was ~16x). Each size keeps the +/// cheapest of three runs, which drops a run that a cache-cold start or a preemption inflated. /// claim: TRN-3 /// defect: D217 #[test] fn a_long_run_of_same_role_messages_merges_in_linear_time() { - let users = vec![json!({"role": "user", "content": "hi"}); 6000]; - let calls = vec![ - json!({"role": "assistant", "content": null, "tool_calls": [ - {"id": "c", "type": "function", "function": {"name": "f", "arguments": "{}"}} - ]}); - 3000 - ]; - let start = std::time::Instant::now(); - let v = c2m( - &chat(json!({"messages": users.into_iter().chain(calls).collect::>()})), - "claude-haiku-4-5", - ); - let took = start.elapsed(); - let blocks = |role: &str| -> usize { + fn thread_cpu() -> std::time::Duration { + let t = rustix::time::clock_gettime(rustix::time::ClockId::ThreadCPUTime); + std::time::Duration::new(t.tv_sec as u64, t.tv_nsec as u32) + } + // `users` user turns, then half as many assistant messages carrying one tool call each. + let request = |users: usize| { + let calls = vec![ + json!({"role": "assistant", "content": null, "tool_calls": [ + {"id": "c", "type": "function", "function": {"name": "f", "arguments": "{}"}} + ]}); + users / 2 + ]; + let user = vec![json!({"role": "user", "content": "hi"}); users]; + chat(json!({"messages": user.into_iter().chain(calls).collect::>()})) + }; + let blocks = |v: &Value, role: &str| -> usize { v["messages"] .as_array() .unwrap() @@ -1722,8 +1728,31 @@ fn a_long_run_of_same_role_messages_merges_in_linear_time() { .map(|m| m["content"].as_array().map_or(1, Vec::len)) .sum() }; - assert_eq!((blocks("user"), blocks("assistant")), (6000, 3000)); - assert!(took < std::time::Duration::from_secs(2), "took {took:?}"); + let cheapest = |users: usize| { + let body = request(users); + (0..3) + .map(|_| { + let start = thread_cpu(); + let v = c2m(&body, "claude-haiku-4-5"); + let took = thread_cpu() - start; + assert_eq!( + (blocks(&v, "user"), blocks(&v, "assistant")), + (users, users / 2) + ); + took + }) + .min() + .unwrap() + }; + // 1,000 user turns cost ~10 ms of CPU in a debug build (far above the clock's resolution); + // the quadratic merge took seconds at 4x that. + let small = cheapest(1_000); + let large = cheapest(4_000); + let ratio = large.as_secs_f64() / small.as_secs_f64(); + assert!( + ratio < 8.0, + "4x the messages cost {ratio:.1}x the CPU ({small:?} -> {large:?}): not linear" + ); } /// A Chat client's `cache_control` on a whole assistant or tool message (the same message-level diff --git a/crates/gateway/tests/claims_security.rs b/crates/gateway/tests/claims_security.rs index 874dba3c..a118cb8f 100644 --- a/crates/gateway/tests/claims_security.rs +++ b/crates/gateway/tests/claims_security.rs @@ -775,30 +775,60 @@ async fn the_rate_limit_counts_one_identity_across_credential_locations() { .await; let vk = vkey(&sk, 12); let url = format!("{}/openai/v1/chat/completions", gw.url()); - // Two requests per location, six in all. Separate buckets would see two each and never trip a - // limit of two; one bucket sees six inside at most two windows, so at least one window has 3. - let mut statuses = Vec::new(); - for round in 0..2 { - for loc in 0..3 { - let mut req = test_client() - .post(if loc == 2 { - format!("{url}?key={vk}") - } else { - url.clone() - }) - .header("content-type", "application/json") - .body(CHAT); - req = match loc { - 0 => req.header("authorization", format!("Bearer {vk}")), - 1 => req.header("x-api-key", vk.clone()), - _ => req, - }; - statuses.push((round, loc, req.send().await.unwrap().status().as_u16())); + // Two requests per location, six in all, sent at once. Separate buckets would see two each and + // never trip a limit of two; one bucket that sees six inside a second (at most two of the + // limiter's whole-second windows) has 3 in at least one. A loaded host can spread the six + // over more windows than that, which proves nothing either way: such a burst is re-sent, as a + // fresh identity so no earlier burst's count carries into it. + let send_six = |vk: String| { + let url = url.clone(); + async move { + let started = Instant::now(); + let mut set = tokio::task::JoinSet::new(); + for round in 0..2 { + for loc in 0..3 { + let mut req = test_client() + .post(if loc == 2 { + format!("{url}?key={vk}") + } else { + url.clone() + }) + .header("content-type", "application/json") + .body(CHAT); + req = match loc { + 0 => req.header("authorization", format!("Bearer {vk}")), + 1 => req.header("x-api-key", vk.clone()), + _ => req, + }; + set.spawn( + async move { (round, loc, req.send().await.unwrap().status().as_u16()) }, + ); + } + } + let mut statuses = set.join_all().await; + statuses.sort(); + (statuses, started.elapsed()) + } + }; + let mut bursts = Vec::new(); + for tenant in 0..20 { + let (statuses, spread) = send_six(if tenant == 0 { + vk.clone() + } else { + vkey(&sk, 1200 + tenant) + }) + .await; + let tripped = statuses.iter().any(|(_, _, s)| *s == 429); + bursts.push((statuses, spread)); + if tripped || spread < Duration::from_secs(1) { + break; } } assert!( - statuses.iter().any(|(_, _, s)| *s == 429), - "one identity in three locations must share one bucket: {statuses:?}" + bursts + .last() + .is_some_and(|(statuses, _)| statuses.iter().any(|(_, _, s)| *s == 429)), + "one identity in three locations must share one bucket: {bursts:?}" ); let respelled = vk.replacen("bai_v1.1.", "bai_v1.01.", 1); @@ -1134,17 +1164,24 @@ async fn a_deny_lands_within_the_bound_and_in_flight_streams_finish() { } assert_eq!(up.hits(), 1, "the stream is in flight before the deny"); + // The two seconds, unless this host is loaded enough to make every request slow (see + // `stretched`): measured now, on the same requests the poll below makes. + let round_trip = median_round_trip(|| async { post().await.status() }).await; + let bound = stretched(Duration::from_secs(2), round_trip); put_kv(nats.port, "blackhole.1717", b"spend").await; let written = Instant::now(); let mut status = 0; - while written.elapsed() < Duration::from_secs(2) { + while written.elapsed() < bound { status = post().await.status().as_u16(); if status == 402 { break; } tokio::time::sleep(Duration::from_millis(50)).await; } - assert_eq!(status, 402, "the deny did not land within 2s"); + assert_eq!( + status, 402, + "the deny did not land within {bound:?} (requests take {round_trip:?})" + ); let (status, body) = in_flight.await.unwrap(); assert_eq!( diff --git a/crates/gateway/tests/claims_streaming.rs b/crates/gateway/tests/claims_streaming.rs index d69d7651..073f245e 100644 --- a/crates/gateway/tests/claims_streaming.rs +++ b/crates/gateway/tests/claims_streaming.rs @@ -10,7 +10,9 @@ mod common; use common::*; use serde_json::Value; -use std::time::{Duration, Instant}; +use std::sync::Arc; +use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering::SeqCst}; +use std::time::Duration; /// How long the scripted provider pauses between its first event and the rest of the stream. const PAUSE: Duration = Duration::from_millis(1500); @@ -49,32 +51,43 @@ async fn pausing_provider(first: &'static str, rest: &'static str) -> ScriptedUp .await } -/// Read a streamed response, noting when `FIRST` arrived and when the stream ended. Runs `on_first` -/// the moment the first event is in hand. -async fn read_stream( - resp: reqwest::Response, - sent: Instant, - mut on_first: impl FnMut(), -) -> (Duration, Duration, String) { +/// A provider that sends `first`, then waits for `gate` before it sends `rest` and closes. +async fn gated_provider( + first: &'static str, + rest: &'static str, + gate: Arc, +) -> ScriptedUpstream { + ScriptedUpstream::start(move |_, _| { + let mut head = http_head(200, "text/event-stream", None); + head.extend_from_slice(first.as_bytes()); + let gate = gate.clone(); + vec![ + Step::Write(head), + Step::Until(Arc::new(move || gate.load(SeqCst))), + Step::Write(rest.as_bytes().to_vec()), + ] + }) + .await +} + +/// Read a streamed response to its end. Runs `on_first` the moment the first event is in hand. +async fn read_stream(resp: reqwest::Response, mut on_first: impl FnMut()) -> String { let mut resp = resp; let mut text = String::new(); - let mut first_at = None; - while let Some(chunk) = tokio::time::timeout(Duration::from_secs(10), resp.chunk()) + let mut seen_first = false; + while let Some(chunk) = tokio::time::timeout(CONDITION_BUDGET, resp.chunk()) .await .expect("the stream must not hang") .unwrap() { text.push_str(&String::from_utf8_lossy(&chunk)); - if first_at.is_none() && text.contains("FIRST") { - first_at = Some(sent.elapsed()); + if !seen_first && text.contains("FIRST") { + seen_first = true; on_first(); } } - ( - first_at.expect("the first event arrived"), - sent.elapsed(), - text, - ) + assert!(seen_first, "the first event arrived: {text}"); + text } async fn post_stream(gw: &Gateway, key: &str, path: &str, body: &str) -> reqwest::Response { @@ -89,7 +102,9 @@ async fn post_stream(gw: &Gateway, key: &str, path: &str, body: &str) -> reqwest } /// The first event reaches the client while the provider is still mid-stream, on a byte relay and -/// on a translated walk alike — the gateway never holds a stream back to the end. +/// on a translated walk alike — the gateway never holds a stream back to the end. The provider +/// sends the rest only once the client has the first event, so the stream completing at all is +/// the proof: an ordering, which no host load can reorder, where a timing would be. /// claim: S1 #[tokio::test] async fn streams_reach_the_client_as_the_provider_sends_them() { @@ -103,7 +118,8 @@ async fn streams_reach_the_client_as_the_provider_sends_them() { ), ] { let (pubkey, sk) = test_keypair(70); - let provider = pausing_provider(first, rest).await; + let gate = Arc::new(AtomicBool::new(false)); + let provider = gated_provider(first, rest, gate.clone()).await; let gw = Gateway::builder(unused_nats_port(), &provider.authority(), &b64(&pubkey)) .providers(providers) .start() @@ -111,19 +127,13 @@ async fn streams_reach_the_client_as_the_provider_sends_them() { let body = format!( r#"{{"model":"{model}","stream":true,"messages":[{{"role":"user","content":"hi"}}]}}"# ); - let sent = Instant::now(); let resp = post_stream(&gw, &billing_vkey(&sk, 70), "/v1/chat/completions", &body).await; assert_eq!(resp.status().as_u16(), 200, "{model}"); - let (first_at, done_at, text) = read_stream(resp, sent, || {}).await; - assert!(text.contains("LAST"), "{model}: {text}"); + let text = read_stream(resp, || gate.store(true, SeqCst)).await; + let (first, last) = (text.find("FIRST"), text.find("LAST")); assert!( - first_at < PAUSE * 2 / 3, - "{model}: the first event took {first_at:?}; the provider sent it at once" - ); - assert!( - done_at - first_at >= PAUSE / 2, - "{model}: the first event must arrive before the provider's pause ends \ - (first {first_at:?}, done {done_at:?})" + first.is_some() && last.is_some() && first < last, + "{model}: {text}" ); } } @@ -131,12 +141,12 @@ async fn streams_reach_the_client_as_the_provider_sends_them() { /// A same-wire Chat Completions stream from a host other than OpenAI goes through `SseBridge`'s /// relay (it drops the `role` OpenRouter repeats on every chunk). Each event still leaves the /// gateway as it arrives: an OpenRouter stream of a Claude row, written one event at a time with -/// gaps and OpenRouter's keep-alive comments between them, reaches the client spread the same way. +/// OpenRouter's keep-alive comments between them, reaches the client one event at a time: the +/// provider writes each event only once the client holds the one before it. /// claim: S1 /// defect: D244 #[tokio::test] async fn a_same_wire_openrouter_chat_relay_streams_each_event_as_it_arrives() { - const GAP: Duration = Duration::from_millis(300); let chunk = |text: &str, finish: &str| { format!( "data: {{\"id\":\"gen-1\",\"object\":\"chat.completion.chunk\",\"created\":1,\ @@ -146,17 +156,25 @@ async fn a_same_wire_openrouter_chat_relay_streams_each_event_as_it_arrives() { ) }; let words = ["W0", "W1", "W2", "W3", "W4"]; + // How many of `words` the client holds. + let seen = Arc::new(AtomicUsize::new(0)); + let client_has = { + let seen = seen.clone(); + move |n: usize| { + let seen = seen.clone(); + Step::Until(Arc::new(move || seen.load(SeqCst) >= n)) + } + }; let provider = ScriptedUpstream::start(move |_, _| { let mut head = http_head(200, "text/event-stream", None); head.extend_from_slice(chunk(words[0], "null").as_bytes()); let mut steps = vec![Step::Write(head)]; - for w in &words[1..] { - steps.push(Step::Sleep(GAP / 2)); + for (i, w) in words.iter().enumerate().skip(1) { + steps.push(client_has(i)); steps.push(Step::Write(b": OPENROUTER PROCESSING\n\n".to_vec())); - steps.push(Step::Sleep(GAP / 2)); steps.push(Step::Write(chunk(w, "null").into_bytes())); } - steps.push(Step::Sleep(GAP)); + steps.push(client_has(words.len())); let mut tail = chunk("", "\"stop\""); tail.push_str( "data: {\"id\":\"gen-1\",\"object\":\"chat.completion.chunk\",\"created\":1,\ @@ -173,61 +191,64 @@ async fn a_same_wire_openrouter_chat_relay_streams_each_event_as_it_arrives() { .start() .await; let body = r#"{"model":"claude-haiku-4-5","stream":true,"stream_options":{"include_usage":true},"messages":[{"role":"user","content":"hi"}]}"#; - let sent = Instant::now(); let mut resp = post_stream(&gw, &billing_vkey(&sk, 72), "/v1/chat/completions", body).await; assert_eq!(resp.status().as_u16(), 200); let mut text = String::new(); - let mut seen_at = Vec::new(); - while let Some(c) = tokio::time::timeout(Duration::from_secs(10), resp.chunk()) + // A word the provider has not written yet cannot be in hand: each one arriving is the gateway + // having relayed the one before it on its own, not held it for the next. + while let Some(c) = tokio::time::timeout(CONDITION_BUDGET, resp.chunk()) .await - .expect("the stream must not hang") + .expect("the stream must not hang: an event was held back for the next") .unwrap() { text.push_str(&String::from_utf8_lossy(&c)); - while seen_at.len() < words.len() && text.contains(words[seen_at.len()]) { - seen_at.push(sent.elapsed()); + let mut n = seen.load(SeqCst); + while n < words.len() && text.contains(words[n]) { + n += 1; } + seen.store(n, SeqCst); } - assert_eq!(seen_at.len(), words.len(), "every word arrived: {text}"); + assert_eq!(seen.load(SeqCst), words.len(), "every word arrived: {text}"); assert!(text.contains("[DONE]"), "{text}"); assert_eq!( text.matches("\"role\"").count(), 1, "the repeated role is dropped: {text}" ); - for (i, pair) in seen_at.windows(2).enumerate() { - assert!( - pair[1] - pair[0] >= GAP / 2, - "{} arrived {:?} after {}; the provider sent them {GAP:?} apart (all: {seen_at:?})", - words[i + 1], - pair[1] - pair[0], - words[i] - ); - } } /// A translated stream whose provider sends only keep-alives while the model thinks (Anthropic's /// `event: ping`, here to a Chat Completions client) still hands its client a byte per ping: an /// SSE comment, which every SSE parser ignores, so a load balancer's idle timeout does not cut a -/// stream that is alive. The answer that follows, and the row, are as without the pings. +/// stream that is alive. The answer that follows, and the row, are as without the pings. Each ping +/// is sent only once the client holds the comment for the one before, so each comment is shown to +/// leave as its ping arrives, whatever the host's load. /// claim: S1 /// claim: REL-1 /// defect: D251 #[tokio::test] async fn a_translated_stream_turns_provider_pings_into_keep_alive_comments() { - const GAP: Duration = Duration::from_millis(250); const PINGS: usize = 4; + // Keep-alive comments the client holds; `usize::MAX` until it holds the first event. + let alive = Arc::new(AtomicUsize::new(usize::MAX)); + let client_has = { + let alive = alive.clone(); + move |n: usize| { + let alive = alive.clone(); + Step::Until(Arc::new(move || alive.load(SeqCst).wrapping_add(1) > n)) + } + }; let provider = ScriptedUpstream::start(move |_, _| { let mut head = http_head(200, "text/event-stream", None); head.extend_from_slice(CLAUDE_FIRST.as_bytes()); let mut steps = vec![Step::Write(head)]; - for _ in 0..PINGS { - steps.push(Step::Sleep(GAP)); + for i in 0..PINGS { + steps.push(client_has(i)); steps.push(Step::Write( b"event: ping\ndata: {\"type\": \"ping\"}\n\n".to_vec(), )); } - steps.push(Step::Sleep(GAP)); + steps.push(client_has(PINGS)); steps.push(Step::Write(CLAUDE_REST.as_bytes().to_vec())); steps }) @@ -241,18 +262,18 @@ async fn a_translated_stream_turns_provider_pings_into_keep_alive_comments() { let mut resp = post_stream(&gw, &billing_vkey(&sk, 73), "/v1/chat/completions", body).await; assert_eq!(resp.status().as_u16(), 200); let mut text = String::new(); - let mut alive_before_last = 0; - while let Some(c) = tokio::time::timeout(Duration::from_secs(10), resp.chunk()) + while let Some(c) = tokio::time::timeout(CONDITION_BUDGET, resp.chunk()) .await - .expect("the stream must not hang") + .expect("the stream must not hang: a ping did not reach the client as a comment") .unwrap() { - let c = String::from_utf8_lossy(&c); - if !text.contains("LAST") && !c.contains("LAST") { - alive_before_last += c.matches(": keep-alive\n\n").count(); + text.push_str(&String::from_utf8_lossy(&c)); + if text.contains("FIRST") { + alive.store(text.matches(": keep-alive\n\n").count(), SeqCst); } - text.push_str(&c); } + let before_last = &text[..text.find("LAST").unwrap_or(text.len())]; + let alive_before_last = before_last.matches(": keep-alive\n\n").count(); assert_eq!( alive_before_last, PINGS, "one keep-alive comment per provider ping, before the answer: {text}" @@ -293,7 +314,7 @@ async fn sigterm_drains_an_open_stream() { ) .await; assert_eq!(resp.status().as_u16(), 200); - let (_, _, text) = read_stream(resp, Instant::now(), || gw.sigterm()).await; + let text = read_stream(resp, || gw.sigterm()).await; assert!( text.contains("LAST") && text.contains("[DONE]"), "the open stream must complete across SIGTERM: {text}" diff --git a/crates/gateway/tests/common/mod.rs b/crates/gateway/tests/common/mod.rs index 27395805..61cc0f76 100644 --- a/crates/gateway/tests/common/mod.rs +++ b/crates/gateway/tests/common/mod.rs @@ -207,6 +207,39 @@ fn subprocess_port_range() -> (u16, u16) { /// this bound while always holding the most recent lines — which is what every assertion reads. const LOG_CAPTURE_CAP: usize = 512 * 1024; +/// A stated latency bound (a deny lands within 2 s), held exactly on an idle host and stretched on +/// a loaded one, where nothing local can meet a wall-clock promise. `round_trip` is +/// [`median_round_trip`] of an ordinary request through the same gateway, measured in the same run: +/// the bound is the larger of the stated one and [`LOAD_STRETCH`] such round trips. On an idle host +/// a round trip is ~35 ms (debug gateway, fresh client connection), so the stated bound is what +/// binds; under a load that makes every hop slow (600 ms round trips at a load average of 160), the +/// bound grows with it. A regression that ignores the event altogether (a deny that never +/// lands) still fails, at whichever bound is in force. +pub fn stretched(stated: Duration, round_trip: Duration) -> Duration { + stated.max(round_trip * LOAD_STRETCH) +} + +/// Round trips a stretched bound may take. A deny lands within one to three round trips at every +/// load measured; 40 leaves an idle host's bound at the stated one (40 x 35 ms is under 2 s), so a +/// deny that lands a second late there still fails. +pub const LOAD_STRETCH: u32 = 40; + +/// The median of five timed runs of `op` (an ordinary request, for [`stretched`]). +pub async fn median_round_trip(mut op: F) -> Duration +where + F: FnMut() -> Fut, + Fut: std::future::Future, +{ + let mut took = Vec::with_capacity(5); + for _ in 0..5 { + let start = std::time::Instant::now(); + let _ = op().await; + took.push(start.elapsed()); + } + took.sort(); + took[2] +} + /// How long a wait on a *condition* (a metric reaching a value, a log line appearing) may take /// before it fails the test. /// @@ -1915,6 +1948,10 @@ impl GatewayBuilder { pub enum Step { Write(Vec), Sleep(Duration), + /// Wait, however long it takes, until the test says so: a reply whose next bytes must follow + /// something the client saw, as an ordering rather than a delay a loaded host could outrun. + /// A gateway that never lets the client see it hangs the test at its own read bound. + Until(Arc bool + Send + Sync>), } /// A complete HTTP/1.1 response with `content-length` and `connection: close`. @@ -2062,6 +2099,11 @@ async fn scripted_conn( let _ = stream.flush().await; } Step::Sleep(d) => sleep(d).await, + Step::Until(ready) => { + while !ready() { + sleep(Duration::from_millis(5)).await; + } + } } } let _ = stream.shutdown().await; diff --git a/crates/gateway/tests/large_bodies.rs b/crates/gateway/tests/large_bodies.rs index 4fb970bf..de94cb3e 100644 --- a/crates/gateway/tests/large_bodies.rs +++ b/crates/gateway/tests/large_bodies.rs @@ -382,9 +382,16 @@ async fn a_large_body_failover_fits_a_tenant_cap_of_one() { .worker_threads(1) .start() .await; + // A held slot shows as a 429; an attempt waiting on the stalled error body, as no answer + // until that upstream's 600s read timeout. Unloaded each takes well under a second, so the + // bound below is a stall guard many times over, not a speed claim a loaded host could fail. + let patient = reqwest::Client::builder() + .timeout(CONDITION_BUDGET * 2) + .build() + .unwrap(); for _ in 0..3 { let started = std::time::Instant::now(); - let resp = client() + let resp = patient .post(format!("{}/v1/chat/completions", gw.url())) .header("authorization", format!("Bearer {}", vkey(&sk))) .header("content-type", "application/json") @@ -394,7 +401,7 @@ async fn a_large_body_failover_fits_a_tenant_cap_of_one() { .unwrap(); assert_eq!(resp.status().as_u16(), 200); assert!( - started.elapsed() < Duration::from_secs(4), + started.elapsed() < CONDITION_BUDGET, "{:?}", started.elapsed() ); diff --git a/crates/gateway/tests/reliability_breaker.rs b/crates/gateway/tests/reliability_breaker.rs index 8cddbd5b..cda00bff 100644 --- a/crates/gateway/tests/reliability_breaker.rs +++ b/crates/gateway/tests/reliability_breaker.rs @@ -195,7 +195,7 @@ async fn a_stalled_half_open_probe_does_not_wedge_the_provider() { // Hold the probe open (unread) while other callers try. let start = Instant::now(); let mut seen = Vec::new(); - while start.elapsed() < Duration::from_secs(8) { + while start.elapsed() < CONDITION_BUDGET { let s = post_byo(&client, &gw.url()).await; seen.push(s); if s == 200 { @@ -207,7 +207,8 @@ async fn a_stalled_half_open_probe_does_not_wedge_the_provider() { assert_eq!( seen.last(), Some(&200), - "with the probe stalled, no request succeeded in 8s ({} tries, upstream hits {}): {seen:?}", + "with the probe stalled, no request succeeded in {CONDITION_BUDGET:?} ({} tries, upstream \ + hits {}): {seen:?}", seen.len(), mock.hits() ); diff --git a/crates/gateway/tests/reliability_early_response.rs b/crates/gateway/tests/reliability_early_response.rs index 9b4a1a53..f9f487d3 100644 --- a/crates/gateway/tests/reliability_early_response.rs +++ b/crates/gateway/tests/reliability_early_response.rs @@ -148,7 +148,9 @@ async fn send_big(gw: &Gateway) -> (u16, String, Duration) { #[tokio::test] async fn an_h2_answer_then_reset_no_error_is_relayed() { let (port, hits) = early_answer_upstream(AfterAnswer::ResetNoError).await; - let gw = gateway(port, 20).await; + // Held behind the body write, the answer would wait out this write timeout: far past the + // bound below, which is in turn many times what relaying it takes unloaded. + let gw = gateway(port, 60).await; let (status, body, took) = send_big(&gw).await; assert_eq!( status, @@ -157,7 +159,7 @@ async fn an_h2_answer_then_reset_no_error_is_relayed() { gw.log() ); assert!(body.contains("maximum context length"), "{body}"); - assert!(took < Duration::from_secs(5), "held for {took:?}"); + assert!(took < Duration::from_secs(30), "held for {took:?}"); assert_eq!( hits.load(Ordering::SeqCst), 1, diff --git a/crates/gateway/tests/reliability_lifecycle.rs b/crates/gateway/tests/reliability_lifecycle.rs index ea4d2164..28ea302f 100644 --- a/crates/gateway/tests/reliability_lifecycle.rs +++ b/crates/gateway/tests/reliability_lifecycle.rs @@ -78,7 +78,10 @@ async fn an_idle_gateway_exits_promptly_on_sigterm() { ); } -/// The same with a short configured grace: shutdown time tracks the grace, even when idle. +/// The same with a short configured grace: shutdown time tracks the grace, even when idle. The +/// drain's own line shows what ended the process (pingora's path out waits the grace, then the +/// runtime timeout), so the time bound need only be the grace, which a loaded host cannot stretch +/// a sub-second drain past. /// claim: REL-16, BIL-21 /// defect: D38 #[tokio::test] @@ -86,16 +89,20 @@ async fn an_idle_gateway_with_a_short_grace_exits_before_the_grace() { let (pubkey, _sk) = test_keypair(1); let mock = MockUpstream::start(Mode::Json).await; let mut gw = Gateway::builder(unused_nats_port(), &mock.authority(), &b64(&pubkey)) - .config_line("shutdown_grace_period_secs = 3") + .config_line("shutdown_grace_period_secs = 10") .config_line("shutdown_runtime_timeout_secs = 1") + // The drain's line is `info`. + .env("AI_LOG", "info") .start() .await; gw.sigterm(); - let exited = gw.wait_exit(Duration::from_secs(15)).await; + let exited = gw.wait_exit(CONDITION_BUDGET).await; assert!( - exited.is_some_and(|t| t < Duration::from_secs(2)), - "idle gateway (3s grace) exit after SIGTERM: {exited:?}" + exited.is_some_and(|t| t < Duration::from_secs(10)), + "idle gateway (10s grace) exit after SIGTERM: {exited:?}" ); + gw.wait_for_log_line(&["drained: no request in flight; exiting"]) + .await; } /// `read_timeout_secs` is honored: a header stall ends at the configured bound. @@ -120,8 +127,10 @@ async fn a_header_stall_ends_at_the_configured_read_timeout() { .map(|r| r.status().as_u16()) .unwrap_or(0); let took = start.elapsed(); + // Unloaded it ends at the 2s bound; the bound it must not wait out instead is the default 600s + // (and the client gives up at 30s, as status 0). assert!( - status >= 500 && took < Duration::from_secs(6), + status >= 500 && took < Duration::from_secs(20), "status {status} after {took:?}" ); } @@ -249,8 +258,9 @@ async fn a_dead_h2_upstream_is_detected_by_ping() { .await; let (status, body, took) = stream_to_end(&gw).await; assert_eq!(status, 200, "{body:?}"); + // Unloaded: about 6s. Under the client's own 40s timeout, which would end it as an error too. assert!( - took < Duration::from_secs(12), + took < Duration::from_secs(30), "a dead upstream was held for {took:?}: {body:?}" ); assert!( @@ -602,7 +612,11 @@ async fn a_panic_in_a_proxy_phase_releases_what_the_request_held() { async fn sigterm_drains_an_in_flight_request_then_exits() { let (pubkey, sk) = test_keypair(1); let mock = MockUpstream::start(Mode::Slow(1500)).await; - let mut gw = Gateway::start(unused_nats_port(), &mock.authority(), &b64(&pubkey)).await; + // `info`: the drain's line, and the billing row. + let mut gw = Gateway::builder(unused_nats_port(), &mock.authority(), &b64(&pubkey)) + .env("AI_LOG", "info") + .start() + .await; let (url, key) = (gw.url(), vkey(&sk, 38)); let held = tokio::spawn(async move { test_client() @@ -616,13 +630,14 @@ async fn sigterm_drains_an_in_flight_request_then_exits() { }); wait_for_metric(&gw, "ai_requests_in_flight", "", 1.0).await; gw.sigterm(); - let exited = gw.wait_exit(Duration::from_secs(15)).await; + // The grace is the default 600s: an exit at all within the budget is the drain's, and its own + // line says so. + let exited = gw.wait_exit(CONDITION_BUDGET).await; let status = held.await.unwrap(); assert_eq!(status.ok(), Some(200), "the in-flight request was cut"); - assert!( - exited.is_some_and(|t| t < Duration::from_secs(5)), - "exit after the drain: {exited:?}" - ); + assert!(exited.is_some(), "no exit after the drain: {exited:?}"); + gw.wait_for_log_line(&["drained: no request in flight; exiting"]) + .await; assert!( gw.log().contains("\"target\":\"ai.usage\""), "the drained request's billing row was not written" diff --git a/crates/gateway/tests/reliability_log_stall.rs b/crates/gateway/tests/reliability_log_stall.rs index 9066b958..820bc24a 100644 --- a/crates/gateway/tests/reliability_log_stall.rs +++ b/crates/gateway/tests/reliability_log_stall.rs @@ -102,11 +102,42 @@ async fn a_rejection_flood_logs_a_capped_number_of_lines() { let resp = client.post(&url).body("{}").send().await.unwrap(); assert_eq!(resp.status(), 404); }; - let flood: u64 = 600; + // Nothing is suppressed unless more than the cap lands in one of the gateway's whole seconds, + // and a loaded host slows the client: flood in rounds until one provably did. `n` rejections + // inside `d` seconds touch at most ceil(d) + 1 of those seconds, so n > 100 * (ceil(d) + 1) + // puts more than 100 into at least one. Eight clients at once, so a round is short. + const ROUND: u64 = 600; + const CLIENTS: u64 = 8; + let mut flood: u64 = 0; let started = std::time::Instant::now(); - for _ in 0..flood { - reject().await; - } + let dense = loop { + let round = std::time::Instant::now(); + let mut set = tokio::task::JoinSet::new(); + for _ in 0..CLIENTS { + let (client, url) = (client.clone(), url.clone()); + set.spawn(async move { + for _ in 0..ROUND / CLIENTS { + let resp = client.post(&url).body("{}").send().await.unwrap(); + assert_eq!(resp.status(), 404); + } + }); + } + while let Some(done) = set.join_next().await { + done.unwrap(); + } + flood += ROUND; + let spread = round.elapsed().as_secs_f64().ceil() as u64; + if ROUND > 100 * (spread + 1) { + break true; + } + if flood >= 20 * ROUND { + break false; + } + }; + assert!( + dense, + "no round of {ROUND} rejections fit in 4s: this host is too loaded to flood the cap" + ); let seconds = started.elapsed().as_secs() + 1; // A fresh second, so this line is admitted and reports what the flood left out. tokio::time::sleep(Duration::from_millis(1100)).await; diff --git a/crates/gateway/tests/replicas.rs b/crates/gateway/tests/replicas.rs index 27fe6abe..2e8eb2fc 100644 --- a/crates/gateway/tests/replicas.rs +++ b/crates/gateway/tests/replicas.rs @@ -259,7 +259,10 @@ async fn a_deny_and_an_allowance_land_on_every_replica_within_the_bound() { for gw in [&a, &b] { assert_eq!(chat(gw, &key).await, 200); } - let bound = Duration::from_secs(2); + // The documented two seconds, unless this host is loaded enough to make every request slow + // (see `stretched`): measured now, on the requests the polls make. + let round_trip = median_round_trip(|| chat(&a, &other)).await; + let bound = stretched(Duration::from_secs(2), round_trip); put_kv(nats.port, "blackhole.4401", b"spend").await; for (gw, name) in [(&a, "A"), (&b, "B")] { @@ -668,14 +671,29 @@ async fn the_per_credential_rate_limit_is_per_replica() { let b = replica().await; let key = billing_vkey(&sk, 4701); + // Bursts of 3 x RPS at once: a burst inside one second spans at most two of the limiter's + // whole-second windows, so one of them holds more than RPS and A must refuse. A loaded host can + // spread a burst wider than that, which proves nothing either way, so it is sent again. let mut refused = false; - for _ in 0..(RPS * 10) { - if chat(&a, &key).await == 429 { - refused = true; + let mut spreads = Vec::new(); + for _ in 0..20 { + let started = Instant::now(); + let mut set = tokio::task::JoinSet::new(); + for _ in 0..RPS * 3 { + let (url, key) = (a.url(), key.clone()); + set.spawn(async move { chat_url(&url, &key).await }); + } + refused = set.join_all().await.contains(&429); + let spread = started.elapsed(); + spreads.push(spread); + if refused || spread < Duration::from_secs(1) { break; } } - assert!(refused, "A refuses the credential past {RPS} req/s"); + assert!( + refused, + "A refuses the credential past {RPS} req/s (bursts took {spreads:?})" + ); // B has seen none of these requests: its window holds 0, so its own RPS are admitted now. for i in 0..RPS { assert_eq!(chat(&b, &key).await, 200, "B request {i} after A's 429"); diff --git a/crates/gateway/tests/request_deadlines.rs b/crates/gateway/tests/request_deadlines.rs index b5ee7516..d17c602c 100644 --- a/crates/gateway/tests/request_deadlines.rs +++ b/crates/gateway/tests/request_deadlines.rs @@ -23,6 +23,24 @@ use std::time::{Duration, Instant}; const CHAT: &str = r#"{"model":"gpt-4o-mini","messages":[{"role":"user","content":"hi"}]}"#; const CHAT_STREAM: &str = r#"{"model":"gpt-4o-mini","stream":true,"messages":[{"role":"user","content":"hi"}]}"#; +/// How soon a request that a deadline ends must have ended. Each one here ends in 1–3 s unloaded, +/// and whatever would end it if the deadline did not (the provider's 60 s script, a body that +/// trickles in for over a minute, the 600 s silence bound, or never) comes far later. So the bound +/// tells the two apart many times over, and a loaded host stretching the 1–3 s cannot fail it; +/// what ended each request is asserted separately, from the answer, the row and the metrics. +const ENDS_WELL_BEFORE: Duration = Duration::from_secs(20); + +/// A request body that takes far longer than [`ENDS_WELL_BEFORE`] to trickle in at +/// [`TRICKLE_GAP`] a byte: trailing whitespace after valid JSON, so a body that does arrive whole +/// still parses. +fn trickled_body() -> String { + format!("{CHAT}{}", " ".repeat(200)) +} + +/// Shorter than the default `client_read_timeout_secs` (60 s), so only `request_max_secs` ends a +/// trickle. +const TRICKLE_GAP: Duration = Duration::from_millis(300); + const OK_JSON: &str = r#"{"id":"chatcmpl-1","object":"chat.completion","model":"gpt-4o-mini","choices":[{"index":0,"message":{"role":"assistant","content":"ok"},"finish_reason":"stop"}],"usage":{"prompt_tokens":5,"completion_tokens":1,"total_tokens":6}}"#; /// An SSE head and then an event every 200 ms for a minute: a stream that never stops moving. @@ -78,7 +96,7 @@ async fn a_stream_that_never_ends_is_cut_at_request_max_secs() { let started = Instant::now(); let resp = post(&gw, "/openai/v1/chat/completions", &key, CHAT_STREAM).await; assert_eq!(resp.status(), 200); - let body = tokio::time::timeout(Duration::from_secs(15), resp.bytes()) + let body = tokio::time::timeout(CONDITION_BUDGET, resp.bytes()) .await .expect("the stream ended"); let took = started.elapsed(); @@ -87,9 +105,9 @@ async fn a_stream_that_never_ends_is_cut_at_request_max_secs() { "a cut stream ended cleanly: {:?}", body.map(|b| String::from_utf8_lossy(&b).into_owned()) ); - // Two seconds, plus up to one tick of the coarse clock that checks a moving stream. + // Never before the two seconds; unloaded, within one tick of the coarse clock after them. assert!( - took >= Duration::from_millis(1500) && took < Duration::from_secs(6), + took >= Duration::from_millis(1500) && took < ENDS_WELL_BEFORE, "{took:?}" ); let row = usage_row_of(&gw).await; @@ -162,7 +180,7 @@ async fn a_silent_provider_is_a_504_at_request_max_secs() { assert_eq!(status, 504, "{path}: {text}\n{}", gw.log()); assert!(text.contains("maximum duration"), "{path}: {text}"); assert!( - took >= Duration::from_millis(1900) && took < Duration::from_secs(5), + took >= Duration::from_millis(1900) && took < ENDS_WELL_BEFORE, "{path}: {took:?}" ); let row = usage_row_of(&gw).await; @@ -273,7 +291,7 @@ async fn h2_stalled_upload(port: u16, headers: &[(&str, String)], sent: &str) -> let (resp, mut send) = client.send_request(req.body(()).unwrap(), false).unwrap(); send.send_data(bytes::Bytes::copy_from_slice(sent.as_bytes()), false) .unwrap(); - let resp = tokio::time::timeout(Duration::from_secs(10), resp) + let resp = tokio::time::timeout(CONDITION_BUDGET, resp) .await .expect("an answer") .unwrap(); @@ -303,7 +321,7 @@ async fn h1_stalled_upload(port: u16, headers: &[(&str, String)], sent: &str) -> s.write_all(head.as_bytes()).await.unwrap(); s.write_all(sent.as_bytes()).await.unwrap(); let mut buf = [0u8; 4096]; - let n = tokio::time::timeout(Duration::from_secs(10), s.read(&mut buf)) + let n = tokio::time::timeout(CONDITION_BUDGET, s.read(&mut buf)) .await .expect("an answer") .unwrap_or(0); @@ -356,7 +374,7 @@ async fn a_client_stalled_in_an_up_front_body_read_gets_a_408() { let took = started.elapsed(); assert_eq!(status, 408, "{what}: {text}\nlog:\n{}", gw.log()); assert!(text.contains("request body timed out"), "{what}: {text}"); - assert!(took < Duration::from_secs(5), "{what}: {took:?}"); + assert!(took < ENDS_WELL_BEFORE, "{what}: {took:?}"); assert_eq!(conns.load(SeqCst), 0, "{what}: reached the upstream"); // The slot is back: under a ceiling of 1 the next request is admitted, which the silent @@ -407,6 +425,7 @@ async fn an_up_front_body_still_trickling_at_request_max_secs_gets_a_408() { .start() .await; let key = billing_vkey(&sk, 268); + let body = trickled_body(); let tcp = tokio::net::TcpStream::connect(("127.0.0.1", gw.port)) .await .unwrap(); @@ -415,27 +434,27 @@ async fn an_up_front_body_still_trickling_at_request_max_secs_gets_a_408() { let req = http::Request::post("http://gw/auto/chat/completions") .header("authorization", format!("Bearer {key}")) .header("content-type", "application/json") - .header("content-length", CHAT.len().to_string()) + .header("content-length", body.len().to_string()) .body(()) .unwrap(); let (resp, mut send) = client.send_request(req, false).unwrap(); let trickle = tokio::spawn(async move { - for b in CHAT.bytes() { + for b in body.bytes() { if send.send_data(bytes::Bytes::from(vec![b]), false).is_err() { return; } - tokio::time::sleep(Duration::from_millis(300)).await; + tokio::time::sleep(TRICKLE_GAP).await; } }); let started = Instant::now(); - let resp = tokio::time::timeout(Duration::from_secs(10), resp) + let resp = tokio::time::timeout(CONDITION_BUDGET, resp) .await .expect("an answer") .unwrap(); let took = started.elapsed(); trickle.abort(); assert_eq!(resp.status().as_u16(), 408, "{}", gw.log()); - assert!(took < Duration::from_secs(5), "{took:?}"); + assert!(took < ENDS_WELL_BEFORE, "{took:?}"); assert_eq!(conns.load(SeqCst), 0, "reached the upstream"); } @@ -470,17 +489,17 @@ async fn a_streamed_upload_still_trickling_at_request_max_secs_ends() { ); wr.write_all(head.as_bytes()).await.unwrap(); let trickle = tokio::spawn(async move { - for b in CHAT.bytes() { + for b in trickled_body().bytes() { let chunk = format!("1\r\n{}\r\n", b as char); if wr.write_all(chunk.as_bytes()).await.is_err() { return; } - tokio::time::sleep(Duration::from_millis(300)).await; + tokio::time::sleep(TRICKLE_GAP).await; } }); let started = Instant::now(); let mut buf = [0u8; 4096]; - let n = tokio::time::timeout(Duration::from_secs(10), rd.read(&mut buf)) + let n = tokio::time::timeout(CONDITION_BUDGET, rd.read(&mut buf)) .await .expect("an answer") .unwrap_or(0); @@ -489,7 +508,7 @@ async fn a_streamed_upload_still_trickling_at_request_max_secs_ends() { let text = String::from_utf8_lossy(&buf[..n]); assert!(text.starts_with("HTTP/1.1 504"), "{text}\n{}", gw.log()); assert!(text.contains("maximum duration"), "{text}"); - assert!(took < Duration::from_secs(5), "{took:?}"); + assert!(took < ENDS_WELL_BEFORE, "{took:?}"); assert_eq!(conns.load(SeqCst), 1); } @@ -521,7 +540,7 @@ async fn a_large_body_re_run_keeps_the_requests_deadline() { let text = resp.text().await.unwrap_or_default(); assert_eq!(status, 504, "{text}\n{}", gw.log()); assert!(text.contains("maximum duration"), "{text}"); - assert!(took < Duration::from_secs(5), "{took:?}"); + assert!(took < ENDS_WELL_BEFORE, "{took:?}"); assert_eq!(conns.load(SeqCst), 1); } From 72c13af58d0431151acef9d09d045384c129b1fd Mon Sep 17 00:00:00 2001 From: Jared Lunde Date: Sun, 4 Oct 2026 12:53:22 -0700 Subject: [PATCH 2/2] gateway: dprint order for the rustix dev-dependency Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01JimHGjsfk2Ktm5GxyZJKKk --- crates/gateway/Cargo.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/gateway/Cargo.toml b/crates/gateway/Cargo.toml index d9e828ce..6904cda1 100644 --- a/crates/gateway/Cargo.toml +++ b/crates/gateway/Cargo.toml @@ -89,11 +89,11 @@ hyper-util = { version = "0.1", features = ["tokio", "server-auto"] } # fork/timeout runner, so no rusty-fork or tempfile in the build graph. proptest = { version = "1.10", default-features = false, features = ["std"] } rcgen = "0.13" +reqwest = { version = "0.13", default-features = false, features = ["json", "rustls", "http2"] } # Thread CPU time (`clock_gettime(CLOCK_THREAD_CPUTIME_ID)`) behind a safe API, for perf assertions # that must hold on a loaded host (`unsafe_code` is forbidden, so not raw `libc`). Already in the # build graph via `tempfile` and `clap`: this only adds its `time` feature. rustix = { version = "1", default-features = false, features = ["std", "time"] } -reqwest = { version = "0.13", default-features = false, features = ["json", "rustls", "http2"] } tokio-rustls = "0.26" [[bench]]