From ddcfa3c5f38816bd9ebeaf52714149204443db0c Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 14 Sep 2026 10:05:18 +0000 Subject: [PATCH] fix(http): avoid delayed-ACK stalls in small relay responses Accepted sockets retained Nagle while the relay flushes individual EVENT and EOSE frames. A delayed-ACK peer could hold small response batches for about 40 ms, including when the immediate peer is a loopback proxy. Enable TCP_NODELAY before handing each accepted socket to Hyper. Reject only that connection if setting the option fails. Protocol framing, proxy configuration and client deadlines remain unchanged. This addresses a measured transport delay, not every reported production timeout. Retain opt-in delayed-ACK and concurrent LMDB read benchmarks for the HTTP/WebSocket path. The controlled delayed-ACK reproduction originally fell from 1.050 s to 15.17 ms for 25 two-event reads. Benchmarks have bounded completion deadlines and validate complete EVENT/EOSE responses; latency assertions remain opt-in because busy CI workers can distort timings. Validation: the original before/after benchmark passed; the rebuilt PR is validated again after restoring its released SDK dependency. --- CHANGELOG.md | 5 ++ docs/explanation/architecture.md | 3 + src/http/mod.rs | 7 ++ tests/relay_response_latency.rs | 146 +++++++++++++++++++++++++++++++ 4 files changed, 161 insertions(+) create mode 100644 tests/relay_response_latency.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index dcbedff..8c42beb 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Fixed + +- Avoid delayed-acknowledgement stalls in small WebSocket and Git responses by + enabling `TCP_NODELAY` on accepted connections. + ## [3.0.2] - 2026-09-11 This release contains no production runtime changes. It improves release and diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index 575900e..f7f359a 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -75,6 +75,9 @@ runtime: - Bind the TCP listener (supports `:0` for kernel-assigned ports; the resolved address back-fills `bind_address` and, if empty, `domain`) +- Disable Nagle on accepted sockets so small EVENT/EOSE and streaming Git + writes do not wait for the peer's delayed acknowledgements, including + when the immediate peer is a local reverse proxy - Initialize Nostr relay builder with custom [`Nip34WritePolicy`](src/nostr/builder.rs:51) - Set up shared storage (LMDB or Memory), purgatory, sync manager, and background maintenance tasks diff --git a/src/http/mod.rs b/src/http/mod.rs index bedd3dc..1c2c2c6 100644 --- a/src/http/mod.rs +++ b/src/http/mod.rs @@ -1073,6 +1073,13 @@ pub async fn run_server_on_listener( loop { let (socket, addr) = listener.accept().await?; + // WebSocket EVENT/EOSE and streaming Git responses flush small writes. + // Nagle can hold the tail until the peer's delayed ACK (about 40 ms + // even through a loopback proxy), multiplying multi-query latency. + if let Err(error) = socket.set_nodelay(true) { + tracing::warn!(peer = %addr, %error, "Could not disable Nagle for client socket"); + continue; + } let io = TokioIo::new(socket); let service = HttpService::new( relay.clone(), diff --git a/tests/relay_response_latency.rs b/tests/relay_response_latency.rs new file mode 100644 index 0000000..52ff71a --- /dev/null +++ b/tests/relay_response_latency.rs @@ -0,0 +1,146 @@ +//! Opt-in Linux delayed-ACK benchmark for the accepted HTTP/WebSocket socket. + +mod common; + +#[tokio::test] +#[ignore = "local load benchmark; run explicitly with --nocapture"] +async fn concurrent_lmdb_read_batches_reach_eose() { + use common::{TestClient, TestRelay}; + use futures_util::{SinkExt, StreamExt}; + use nostr_sdk::prelude::*; + use std::time::{Duration, Instant}; + use tokio_tungstenite::tungstenite::Message; + + let relay = TestRelay::start_with_lmdb().await; + let keys = relay.owner_keys().clone(); + let client = TestClient::new(relay.url(), keys.clone()).await.unwrap(); + for index in 0..32 { + let event = EventBuilder::new(Kind::TextNote, format!("{index}:{}", "x".repeat(8192))) + .finalize(&keys) + .unwrap(); + client.send_event(&event).await.unwrap(); + } + for concurrency in [1, 4, 16, 32] { + let url = relay.url().to_string(); + let reads = (0..concurrency).map(|_| { + let url = url.clone(); + async move { + let (mut stream, _) = tokio_tungstenite::connect_async(url).await.unwrap(); + let started = Instant::now(); + for _ in 0..3 { + stream + .send(Message::Text(r#"["REQ","load",{"kinds":[1]}]"#.into())) + .await + .unwrap(); + let mut events = 0; + loop { + let frame = stream.next().await.unwrap().unwrap(); + if !frame.is_text() { + continue; + } + let message: serde_json::Value = + serde_json::from_str(frame.to_text().unwrap()).unwrap(); + match message[0].as_str() { + Some("EVENT") => events += 1, + Some("EOSE") => break, + other => panic!("unexpected reply: {other:?}"), + } + } + assert_eq!(events, 32); + } + let elapsed = started.elapsed(); + stream.close(None).await.unwrap(); + elapsed + } + }); + let times = tokio::time::timeout( + Duration::from_secs(30), + futures_util::future::join_all(reads), + ) + .await + .expect("all readers must reach EOSE under local load"); + eprintln!( + "{concurrency} concurrent LMDB readers, three 32-event batches each: max {:?}", + times.iter().max().unwrap() + ); + } + client.disconnect().await; + relay.stop().await; +} + +#[cfg(target_os = "linux")] +#[tokio::test] +#[ignore = "socket latency benchmark; run explicitly on an otherwise idle worker"] +async fn small_event_batches_do_not_wait_for_delayed_ack() { + use std::os::fd::AsRawFd; + use std::time::{Duration, Instant}; + + use common::{TestClient, TestRelay}; + use futures_util::{SinkExt, StreamExt}; + use nostr_sdk::prelude::*; + use tokio_tungstenite::tungstenite::Message; + + let relay = TestRelay::start().await; + let keys = relay.owner_keys().clone(); + let client = TestClient::new(relay.url(), keys.clone()).await.unwrap(); + for index in 0..2 { + let event = EventBuilder::new(Kind::TextNote, format!("event {index}")) + .finalize(&keys) + .unwrap(); + client.send_event(&event).await.unwrap(); + } + let socket = tokio::net::TcpStream::connect(relay.domain()) + .await + .unwrap(); + socket.set_nodelay(true).unwrap(); + let (mut stream, _) = tokio_tungstenite::client_async(relay.url(), socket) + .await + .unwrap(); + let started = Instant::now(); + tokio::time::timeout(Duration::from_secs(10), async { + for _ in 0..25 { + let disabled: libc::c_int = 0; + // Explicitly model a peer that uses delayed acknowledgements. + // The descriptor and option value remain alive for setsockopt. + let result = unsafe { + libc::setsockopt( + stream.get_ref().as_raw_fd(), + libc::IPPROTO_TCP, + libc::TCP_QUICKACK, + &disabled as *const _ as *const libc::c_void, + std::mem::size_of_val(&disabled) as libc::socklen_t, + ) + }; + assert_eq!(result, 0); + stream + .send(Message::Text(r#"["REQ","latency",{"kinds":[1]}]"#.into())) + .await + .unwrap(); + let mut events = 0; + loop { + let frame = stream.next().await.unwrap().unwrap(); + if !frame.is_text() { + continue; + } + let message: serde_json::Value = + serde_json::from_str(frame.to_text().unwrap()).unwrap(); + match message[0].as_str() { + Some("EVENT") => events += 1, + Some("EOSE") => break, + other => panic!("unexpected reply: {other:?}"), + } + } + assert_eq!(events, 2); + } + }) + .await + .expect("batch reads must complete"); + let elapsed = started.elapsed(); + eprintln!("25 two-event reads with delayed ACK: {elapsed:?}"); + assert!( + elapsed < Duration::from_millis(500), + "small batches were delayed: {elapsed:?}" + ); + client.disconnect().await; + relay.stop().await; +}