Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Generated
+1
@@ -2571,6 +2571,7 @@ dependencies = [
|
|||||||
"rmpv",
|
"rmpv",
|
||||||
"serde",
|
"serde",
|
||||||
"serde_json",
|
"serde_json",
|
||||||
|
"socket2 0.6.4",
|
||||||
"thiserror",
|
"thiserror",
|
||||||
"tokio",
|
"tokio",
|
||||||
"tracing",
|
"tracing",
|
||||||
|
|||||||
@@ -21,6 +21,7 @@ futures = "0.3"
|
|||||||
pyo3 = { version = "0.29.0", features = ["extension-module"] }
|
pyo3 = { version = "0.29.0", features = ["extension-module"] }
|
||||||
serde = { version = "1", features = ["derive"] }
|
serde = { version = "1", features = ["derive"] }
|
||||||
serde_json = "1"
|
serde_json = "1"
|
||||||
|
socket2 = "0.6"
|
||||||
thiserror = "2"
|
thiserror = "2"
|
||||||
tokio = { version = "1", features = ["full"] }
|
tokio = { version = "1", features = ["full"] }
|
||||||
tokio-stream = { version = "0.1", features = ["net"] }
|
tokio-stream = { version = "0.1", features = ["net"] }
|
||||||
|
|||||||
@@ -24,6 +24,7 @@ futures = { workspace = true }
|
|||||||
pyo3 = { workspace = true }
|
pyo3 = { workspace = true }
|
||||||
serde = { workspace = true }
|
serde = { workspace = true }
|
||||||
serde_json = { workspace = true }
|
serde_json = { workspace = true }
|
||||||
|
socket2 = { workspace = true }
|
||||||
thiserror = { workspace = true }
|
thiserror = { workspace = true }
|
||||||
tokio = { workspace = true }
|
tokio = { workspace = true }
|
||||||
tracing = { workspace = true }
|
tracing = { workspace = true }
|
||||||
|
|||||||
@@ -88,6 +88,14 @@ pub async fn serve(
|
|||||||
// runtime drops → detached handlers cancel → their `AbortGuard`s fire, release
|
// runtime drops → detached handlers cancel → their `AbortGuard`s fire, release
|
||||||
// `Senders` clones → tok/detok channels close → workers exit. Full drain is
|
// `Senders` clones → tok/detok channels close → workers exit. Full drain is
|
||||||
// deferred (see `request_shutdown`).
|
// 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.
|
// `with_connect_info` exposes the peer address to the access-log middleware.
|
||||||
let serve = axum::serve(
|
let serve = axum::serve(
|
||||||
listener,
|
listener,
|
||||||
|
|||||||
@@ -93,9 +93,28 @@ impl Drop for Runtime {
|
|||||||
/// startup misconfiguration (e.g. no tokenizer for a non-skip server).
|
/// startup misconfiguration (e.g. no tokenizer for a non-skip server).
|
||||||
pub fn start(cfg: RuntimeConfig) -> Result<Runtime, String> {
|
pub fn start(cfg: RuntimeConfig) -> Result<Runtime, String> {
|
||||||
// Bind the API server port before spawning any thread, so an unavailable
|
// Bind the API server port before spawning any thread, so an unavailable
|
||||||
// port (EADDRINUSE) is a hard startup error.
|
// port (EADDRINUSE) is a hard startup error. socket2 rather than
|
||||||
let listener = std::net::TcpListener::bind(cfg.rust_server_args.http_addr)
|
// `std::net::TcpListener` so SO_RCVBUF can be set before `listen`.
|
||||||
.map_err(|e| format!("bind {} failed: {e}", cfg.rust_server_args.http_addr))?;
|
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
|
listener
|
||||||
.set_nonblocking(true)
|
.set_nonblocking(true)
|
||||||
.map_err(|e| format!("listener set_nonblocking failed: {e}"))?;
|
.map_err(|e| format!("listener set_nonblocking failed: {e}"))?;
|
||||||
|
|||||||
Reference in New Issue
Block a user