mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
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.
This commit is contained in:
@@ -7,6 +7,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
|||||||
|
|
||||||
## [Unreleased]
|
## [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
|
## [3.0.2] - 2026-09-11
|
||||||
|
|
||||||
This release contains no production runtime changes. It improves release and
|
This release contains no production runtime changes. It improves release and
|
||||||
|
|||||||
@@ -75,6 +75,9 @@ runtime:
|
|||||||
|
|
||||||
- Bind the TCP listener (supports `:0` for kernel-assigned ports; the
|
- Bind the TCP listener (supports `:0` for kernel-assigned ports; the
|
||||||
resolved address back-fills `bind_address` and, if empty, `domain`)
|
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)
|
- Initialize Nostr relay builder with custom [`Nip34WritePolicy`](src/nostr/builder.rs:51)
|
||||||
- Set up shared storage (LMDB or Memory), purgatory, sync manager, and
|
- Set up shared storage (LMDB or Memory), purgatory, sync manager, and
|
||||||
background maintenance tasks
|
background maintenance tasks
|
||||||
|
|||||||
@@ -1073,6 +1073,13 @@ pub async fn run_server_on_listener(
|
|||||||
|
|
||||||
loop {
|
loop {
|
||||||
let (socket, addr) = listener.accept().await?;
|
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 io = TokioIo::new(socket);
|
||||||
let service = HttpService::new(
|
let service = HttpService::new(
|
||||||
relay.clone(),
|
relay.clone(),
|
||||||
|
|||||||
@@ -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;
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user