diff --git a/Cargo.lock b/Cargo.lock index bd250b8af591..b266a730591b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7779,9 +7779,9 @@ dependencies = [ [[package]] name = "quinn" -version = "0.11.9" +version = "0.11.12" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b9e20a958963c291dc322d98411f541009df2ced7b5a4f2bd52337638cfccf20" +checksum = "4051e23e9185c255a7e33ef59cdbca87a22d359052eecd22fc6b901fb37d9d11" dependencies = [ "bytes", "cfg_aliases", @@ -7799,9 +7799,9 @@ dependencies = [ [[package]] name = "quinn-proto" -version = "0.11.17" +version = "0.11.18" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "04759210543be93709136e28212294a659ef5001836ff4eab4d663e4529bba83" +checksum = "a9746dbde176634f4f2f1faf2404e30a31b2bc1e9cafb5329c95d8177a18c9fc" dependencies = [ "bytes", "getrandom 0.4.2", diff --git a/Cargo.toml b/Cargo.toml index e76399a77681..43f114bb913d 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -203,7 +203,7 @@ vllm-proto = { version = "=0.1.0" } kube = { version = "2.0.1", default-features = false, features = ["client", "rustls-tls", "aws-lc-rs", "runtime"] } k8s-openapi = { version = "0.26.0", features = ["v1_32"] } rcgen = { version = "0.13" } -quinn = { version = "0.11.9", default-features = false, features = ["runtime-tokio", "rustls-ring"] } +quinn = { version = "0.11.12", default-features = false, features = ["runtime-tokio", "rustls-ring"] } sha2 = { version = "0.10" } tokio-rustls = { version = "0.26" } rustls = { version = "0.23", default-features = false, features = ["ring", "std"] } diff --git a/lib/bindings/kvbm/Cargo.lock b/lib/bindings/kvbm/Cargo.lock index bbb8ce02a133..094ca612b2b6 100644 --- a/lib/bindings/kvbm/Cargo.lock +++ b/lib/bindings/kvbm/Cargo.lock @@ -6276,9 +6276,9 @@ dependencies = [ [[package]] name = "quinn" -version = "0.11.9" +version = "0.11.12" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b9e20a958963c291dc322d98411f541009df2ced7b5a4f2bd52337638cfccf20" +checksum = "4051e23e9185c255a7e33ef59cdbca87a22d359052eecd22fc6b901fb37d9d11" dependencies = [ "bytes", "cfg_aliases", @@ -6296,9 +6296,9 @@ dependencies = [ [[package]] name = "quinn-proto" -version = "0.11.17" +version = "0.11.18" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "04759210543be93709136e28212294a659ef5001836ff4eab4d663e4529bba83" +checksum = "a9746dbde176634f4f2f1faf2404e30a31b2bc1e9cafb5329c95d8177a18c9fc" dependencies = [ "bytes", "getrandom 0.4.2", diff --git a/lib/bindings/python/Cargo.lock b/lib/bindings/python/Cargo.lock index fcf01abe3d0b..fd8c1fcbc80e 100644 --- a/lib/bindings/python/Cargo.lock +++ b/lib/bindings/python/Cargo.lock @@ -6756,9 +6756,9 @@ dependencies = [ [[package]] name = "quinn" -version = "0.11.9" +version = "0.11.12" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b9e20a958963c291dc322d98411f541009df2ced7b5a4f2bd52337638cfccf20" +checksum = "4051e23e9185c255a7e33ef59cdbca87a22d359052eecd22fc6b901fb37d9d11" dependencies = [ "bytes", "cfg_aliases", @@ -6776,9 +6776,9 @@ dependencies = [ [[package]] name = "quinn-proto" -version = "0.11.17" +version = "0.11.18" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "04759210543be93709136e28212294a659ef5001836ff4eab4d663e4529bba83" +checksum = "a9746dbde176634f4f2f1faf2404e30a31b2bc1e9cafb5329c95d8177a18c9fc" dependencies = [ "bytes", "getrandom 0.4.2", diff --git a/lib/runtime/examples/Cargo.lock b/lib/runtime/examples/Cargo.lock index a55cad92c34e..23f376c5d65b 100644 --- a/lib/runtime/examples/Cargo.lock +++ b/lib/runtime/examples/Cargo.lock @@ -2607,9 +2607,9 @@ dependencies = [ [[package]] name = "quinn" -version = "0.11.9" +version = "0.11.12" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b9e20a958963c291dc322d98411f541009df2ced7b5a4f2bd52337638cfccf20" +checksum = "4051e23e9185c255a7e33ef59cdbca87a22d359052eecd22fc6b901fb37d9d11" dependencies = [ "bytes", "cfg_aliases", @@ -2627,9 +2627,9 @@ dependencies = [ [[package]] name = "quinn-proto" -version = "0.11.17" +version = "0.11.18" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "04759210543be93709136e28212294a659ef5001836ff4eab4d663e4529bba83" +checksum = "a9746dbde176634f4f2f1faf2404e30a31b2bc1e9cafb5329c95d8177a18c9fc" dependencies = [ "bytes", "getrandom 0.4.2", diff --git a/lib/runtime/src/pipeline/network/quic_response.rs b/lib/runtime/src/pipeline/network/quic_response.rs index 2fa95a4cc6f0..a7c31035917c 100644 --- a/lib/runtime/src/pipeline/network/quic_response.rs +++ b/lib/runtime/src/pipeline/network/quic_response.rs @@ -70,11 +70,13 @@ const FRAME_HEADER_LEN: usize = 1 + 16 + 4; // at the receiver so parsing those frames does not poll Quinn once per header // and payload and exhaust Tokio's cooperative task budget. const RECEIVE_BUFFER_CAPACITY: usize = 256 * 1024; +const READ_BUDGET_FRAMES: usize = 16; +const READ_BUDGET_BYTES: usize = 4 * 1024; // A single frontend UDP socket is limited by the host receive-buffer ceiling. // Reuse-port endpoints preserve one advertised address while spreading QUIC // connections and receive queues across several sockets. #[cfg(target_os = "linux")] -const SERVER_ENDPOINTS: usize = 8; +const SERVER_ENDPOINTS: usize = 32; #[cfg(not(target_os = "linux"))] const SERVER_ENDPOINTS: usize = 1; const MAX_FRAME_PAYLOAD: usize = 32 * 1024 * 1024; @@ -854,7 +856,7 @@ fn bind_server_udp(address: SocketAddr, join_reuseport: bool) -> std::io::Result socket.bind(&address.into())?; // Bind the first endpoint exclusively so an ephemeral port cannot join an // unrelated server's reuse-port group. Linux permits enabling reuse-port - // after that first bind, and the remaining seven endpoints can then join. + // after that first bind, and the remaining endpoints can then join. #[cfg(target_os = "linux")] if !join_reuseport { socket.set_reuse_address(true)?; @@ -932,8 +934,9 @@ async fn run_server_connection( ) .await; if let Err(error) = result { - tracing::warn!(connection_id, %error, "QUIC response lane failed; closing connection"); + tracing::warn!(connection_id, error = format!("{error:#}"), close_reason = ?lane_connection.close_reason(), "QUIC response lane failed; closing connection"); fail_server_connection(&lane_state, connection_id); + // Close explicitly if the lane failed before bundle registration. lane_connection .close(CLOSE_CODE_INVARIANT, b"response lane invariant failure"); } @@ -1041,17 +1044,30 @@ async fn run_server_lane( register_server_connection_bundle(&state, connection, bundle_id)?; let (control_tx, mut control_rx) = mpsc::channel::(RESPONSE_BUFFER_CAPACITY); - let mut writer = tokio::spawn(async move { + let writer = async move { while let Some(frame) = control_rx.recv().await { let mut chunks = [frame.header(), frame.payload]; send.write_all_chunks(&mut chunks).await?; } Ok::<(), quinn::WriteError>(()) - }); + }; let reader = async { + let mut frames = READ_BUDGET_FRAMES; + let mut bytes = READ_BUDGET_BYTES; loop { + // Charge a bounded batch of small frames to the task budget. + // Charging every token frame can leave already-received responses + // queued behind repeated scheduler waits. Keep both bounds so + // large frames still yield and the reverse-control writer runs. + if frames >= READ_BUDGET_FRAMES || bytes >= READ_BUDGET_BYTES { + tokio::task::consume_budget().await; + frames = 0; + bytes = 0; + } let frame = read_frame(&mut recv).await?; + frames += 1; + bytes += FRAME_HEADER_LEN + frame.payload.len(); process_server_frame( frame, bundle_id, @@ -1064,16 +1080,13 @@ async fn run_server_lane( #[allow(unreachable_code)] Ok::<(), anyhow::Error>(()) }; - tokio::pin!(reader); + // Keep both futures pinned across polls. This avoids a separate task wake + // for registration acknowledgements without cancelling partial reads. + tokio::pin!(reader, writer); tokio::select! { - result = &mut reader => { - writer.abort(); - let _ = writer.await; - result - } + result = &mut reader => result, result = &mut writer => match result { - Ok(Ok(())) => bail!("QUIC reverse-control writer exited unexpectedly"), - Ok(Err(error)) => Err(error.into()), + Ok(()) => bail!("QUIC reverse-control writer exited unexpectedly"), Err(error) => Err(error.into()), }, } @@ -1848,14 +1861,23 @@ fn spawn_client_lane( &writer_connections, &writer_contexts, &writer_healthy, - &error.to_string(), + bundle_id, + "writer", + &format!("{error:#}"), ); } }); tokio::spawn(async move { if let Err(error) = run_client_control_reader(recv, contexts.clone()).await { - fail_client_connection_bundle(&connections, &contexts, &healthy, &error.to_string()); + fail_client_connection_bundle( + &connections, + &contexts, + &healthy, + bundle_id, + "control_reader", + &format!("{error:#}"), + ); } }); } @@ -1946,13 +1968,22 @@ fn fail_client_connection_bundle( connections: &[quinn::Connection], contexts: &Mutex>>, healthy: &AtomicBool, + bundle_id: Uuid, + failure_path: &'static str, reason: &str, ) { if !healthy.swap(false, Ordering::AcqRel) { return; } crate::metrics::quic_response::record_bundle_failure("worker"); - tracing::warn!(%reason, "QUIC response connection bundle invariant failed"); + let remote = connections.first().map(quinn::Connection::remote_address); + tracing::warn!( + %bundle_id, + failure_path, + ?remote, + reason, + "QUIC response connection bundle invariant failed" + ); for (_, entry) in contexts.lock().drain() { entry.fail_registration(reason); entry.record_cancellation();