From 4494fb96b2fd4933e101cdf44dcc32d974611e5a Mon Sep 17 00:00:00 2001 From: Rain Jiang <96632942+rainj-me@users.noreply.github.com> Date: Mon, 3 Aug 2026 22:10:52 -0700 Subject: [PATCH] refactor the tcp listener binding logic (#33420) --- rust/sglang-server/src/api_server.rs | 17 ------------- rust/sglang-server/src/runtime.rs | 34 ++++++-------------------- rust/sglang-server/src/utils.rs | 1 + rust/sglang-server/src/utils/sock.rs | 36 ++++++++++++++++++++++++++++ 4 files changed, 44 insertions(+), 44 deletions(-) create mode 100644 rust/sglang-server/src/utils/sock.rs diff --git a/rust/sglang-server/src/api_server.rs b/rust/sglang-server/src/api_server.rs index a827275f1..ff3b2cf4e 100644 --- a/rust/sglang-server/src/api_server.rs +++ b/rust/sglang-server/src/api_server.rs @@ -82,23 +82,6 @@ pub async fn serve( return; } }; - if let Ok(addr) = listener.local_addr() { - tracing::info!(%addr, "sglang-server api listening"); - } - // Non-graceful shutdown: on the signal, stop accepting and RETURN without - // waiting for in-flight handlers (a `/generate` blocked on egress would wedge - // the join). Returning unwinds `block_on` in `runtime::start` → the api tokio - // runtime drops → detached handlers cancel → their `AbortGuard`s fire, release - // `Senders` clones → tok/detok channels close → workers exit. Full drain is - // deferred (see `request_shutdown`). - // Match Python (asyncio sets TCP_NODELAY); avoids a ~13 ms - // Nagle/delayed-ACK penalty on keep-alive connections. - use axum::serve::ListenerExt; - let listener = listener.tap_io(|io| { - if let Err(e) = io.set_nodelay(true) { - tracing::debug!(error = %e, "set_nodelay failed"); - } - }); // `with_connect_info` exposes the peer address to the access-log middleware. let serve = axum::serve( listener, diff --git a/rust/sglang-server/src/runtime.rs b/rust/sglang-server/src/runtime.rs index 89c93b886..df5f425d4 100644 --- a/rust/sglang-server/src/runtime.rs +++ b/rust/sglang-server/src/runtime.rs @@ -27,6 +27,7 @@ use crate::ring::{ }; use crate::runtime::threads::{plan_cores, spawn_pool}; use crate::tokenizer_manager::{Senders, TmEvent}; +use crate::utils::sock::bind_tcp_listener; use crate::{api_server, detokenizer, tokenizer, tokenizer_manager}; // Re-export so stages keep importing `crate::runtime::Runnable`. @@ -92,33 +93,6 @@ impl Drop for Runtime { /// so the Python caller regains control of the GIL immediately. `Err` on a /// startup misconfiguration (e.g. no tokenizer for a non-skip server). pub fn start(cfg: RuntimeConfig) -> Result { - // Bind the API server port before spawning any thread, so an unavailable - // port (EADDRINUSE) is a hard startup error. socket2 rather than - // `std::net::TcpListener` so SO_RCVBUF can be set before `listen`. - let addr = cfg.rust_server_args.http_addr; - let socket = socket2::Socket::new( - socket2::Domain::for_address(addr), - socket2::Type::STREAM, - Some(socket2::Protocol::TCP), - ) - .map_err(|e| format!("socket for {addr} failed: {e}"))?; - socket - .set_reuse_address(true) - .map_err(|e| format!("set_reuseaddr failed: {e}"))?; - if let Err(e) = socket.set_recv_buffer_size(16 * 1024 * 1024) { - eprintln!("warning: set_recv_buffer_size on listener failed: {e}"); - } - socket - .bind(&addr.into()) - .map_err(|e| format!("bind {addr} failed: {e}"))?; - socket - .listen(1024) - .map_err(|e| format!("listen on {addr} failed: {e}"))?; - let listener: std::net::TcpListener = socket.into(); - listener - .set_nonblocking(true) - .map_err(|e| format!("listener set_nonblocking failed: {e}"))?; - let (shutdown_tx, shutdown_rx) = flume::unbounded::<()>(); let mut threads = Vec::new(); let plan = plan_cores(&cfg); @@ -268,6 +242,12 @@ pub fn start(cfg: RuntimeConfig) -> Result { let senders = senders.clone(); let api_activity = egress_activity.clone(); let shutdown_rx = shutdown_rx.clone(); + // Bind synchronously so an unavailable port (EADDRINUSE) is a hard + // startup error. The `?` drops `shutdown_tx`/`senders`, which stops the + // launcher process. + let http_addr = cfg.rust_server_args.http_addr; + let listener = bind_tcp_listener(http_addr) + .map_err(|e| format!("binding API listener on {} failed: {e}", http_addr))?; let handle = std::thread::Builder::new() .name("api-runtime".into()) .spawn(move || { diff --git a/rust/sglang-server/src/utils.rs b/rust/sglang-server/src/utils.rs index 2301fcb1d..39d9e0a1e 100644 --- a/rust/sglang-server/src/utils.rs +++ b/rust/sglang-server/src/utils.rs @@ -1,3 +1,4 @@ //! Shared helpers with no home in a pipeline stage. pub mod regex; +pub mod sock; diff --git a/rust/sglang-server/src/utils/sock.rs b/rust/sglang-server/src/utils/sock.rs new file mode 100644 index 000000000..3506a0bd4 --- /dev/null +++ b/rust/sglang-server/src/utils/sock.rs @@ -0,0 +1,36 @@ +//! Socket helpers for the API listener. + +use std::io; +use std::net::SocketAddr; +use std::net::TcpListener; + +const BACKLOG: i32 = 2048; // matches uvicorn's default (asyncio's own is 100) +const RECV_BUF_SIZE: usize = 16 * 1024 * 1024; + +/// Bind and tune the API listener, returning it ready for +/// `tokio::net::TcpListener::from_std`. +/// +/// socket2 rather than `TcpListener` so options (SO_RCVBUF, ...) can +/// be set before `listen`. +pub fn bind_tcp_listener(addr: SocketAddr) -> io::Result { + let socket = socket2::Socket::new( + socket2::Domain::for_address(addr), + socket2::Type::STREAM, + Some(socket2::Protocol::TCP), + )?; + socket.set_reuse_address(true)?; + if let Err(e) = socket.set_recv_buffer_size(RECV_BUF_SIZE) { + tracing::warn!( + "set_recv_buffer_size({RECV_BUF_SIZE}) failed: {e}; continuing with the default size" + ); + } + // Matches Python, accepted sockets inherit TCP_NODELAY. + if let Err(e) = socket.set_tcp_nodelay(true) { + tracing::warn!("set_tcp_nodelay failed: {e}; continuing without TCP_NODELAY"); + } + socket.bind(&addr.into())?; + socket.listen(BACKLOG)?; + let listener: std::net::TcpListener = socket.into(); + listener.set_nonblocking(true)?; + Ok(listener) +}