perf(quartz): TCP_NODELAY for every relay websocket client — kills CLOSE→REQ Nagle stalls

Found while attributing the small-REQ wire floor (backlog item 6,
latency half): geode's new WireReqFloorBenchmark measured a flat
43.7 ms per REQ round trip that survived every server-side change —
store configs, dispatchers, the pump — and then vanished when the
round's preceding CLOSE was dropped. Root cause is client-side: OkHttp
does not set TCP_NODELAY, relays never answer a CLOSE (NIP-01), so its
bytes sit unACKed for the peer's ~40 ms delayed-ACK window and Nagle
holds the next REQ behind them. CLOSE-then-REQ is a Nostr client's
hottest pattern — every feed/filter switch.

relayBench's harness client already shipped a no-delay socket factory
(which is why benchmark numbers never showed the stall) but the
production clients did not. New TcpNoDelaySocketFactory (quartz
jvmAndroid, next to BasicOkHttpWebSocket) is now used by the Android
relay pool factory, the Desktop relay client, amy's relay connections,
and geode's mirror worker. Direct connections only — SOCKS/Tor paths
are untouched.

With the factory, the benchmark puts geode's ~21-row REQ at ~1.25 ms
on the wire (matching relayBench): ~0.6 ms Ktor CIO+OkHttp loopback
floor, ~0.5 ms per-REQ server work (already investigated). Per-frame
burst cost measured negligible and the pump adds ~nothing, so the
send-path latency angle of backlog item 6 is closed as not-a-problem;
its ingest-CPU share remains a separate throughput question. Findings
recorded in quartz/plans/2026-07-04-small-req-floor.md.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01TtDNpayEYvJH7QuPswND3A
This commit is contained in:
Claude
2026-07-04 00:38:09 +00:00
parent f22657c870
commit fd6662ca30
9 changed files with 457 additions and 15 deletions
@@ -20,6 +20,7 @@
*/
package com.vitorpamplona.amethyst.service.okhttp
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.TcpNoDelaySocketFactory
import com.vitorpamplona.quartz.utils.Log
import okhttp3.Dispatcher
import okhttp3.OkHttpClient
@@ -57,6 +58,12 @@ class OkHttpClientFactoryForRelays(
OkHttpClient
.Builder()
.dispatcher(myDispatcher)
// TCP_NODELAY: a CLOSE (never answered by relays) followed by a
// REQ — every feed switch — otherwise nagles the REQ behind the
// unACKed CLOSE for the peer's delayed-ACK window (~40 ms+).
// See quartz TcpNoDelaySocketFactory. Direct connections only;
// the Tor SOCKS path is unaffected.
.socketFactory(TcpNoDelaySocketFactory)
.dns(dns)
.eventListenerFactory(DnsInvalidatingEventListener.Factory(dns))
.followRedirects(true)
@@ -50,6 +50,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.BasicOkHttpWebSocket
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.TcpNoDelaySocketFactory
import com.vitorpamplona.quartz.nip01Core.signers.NostrSigner
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal
import com.vitorpamplona.quartz.nip01Core.store.IEventStore
@@ -111,7 +112,7 @@ class Context(
val identity: Identity,
val state: RunState,
) : AutoCloseable {
private val okhttp = OkHttpClient.Builder().build()
private val okhttp = OkHttpClient.Builder().socketFactory(TcpNoDelaySocketFactory).build()
val client: NostrClient =
NostrClient(
@@ -31,6 +31,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.client.single.newSubId
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.BasicOkHttpWebSocket
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.TcpNoDelaySocketFactory
import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.channels.Channel.Factory.UNLIMITED
import kotlinx.coroutines.selects.select
@@ -143,7 +144,7 @@ object NipCommand {
timeoutMs: Long,
): List<Map<String, Any?>> {
if (SEARCH_RELAYS.isEmpty()) return emptyList()
val okhttp = OkHttpClient.Builder().build()
val okhttp = OkHttpClient.Builder().socketFactory(TcpNoDelaySocketFactory).build()
val client = NostrClient(websocketBuilder = BasicOkHttpWebSocket.Builder { okhttp })
val filter = Filter(kinds = listOf(NIPTEXT_KIND, WIKI_KIND, LONGFORM_KIND), search = "NIP-$slug", limit = 10)
val subId = newSubId()
@@ -35,6 +35,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.BasicOkHttpWebSocket
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.TcpNoDelaySocketFactory
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal
import com.vitorpamplona.quartz.nip46RemoteSigner.BunkerResponse
import com.vitorpamplona.quartz.nip46RemoteSigner.NostrConnectEvent
@@ -128,7 +129,7 @@ object NostrConnect {
System.err.println("[nostrconnect] paste this into your signer within ${timeoutMs / 1000}s:")
System.err.println(offer)
val okhttp = OkHttpClient.Builder().build()
val okhttp = OkHttpClient.Builder().socketFactory(TcpNoDelaySocketFactory).build()
val client = NostrClient(websocketBuilder = BasicOkHttpWebSocket.Builder { okhttp })
val incoming = Channel<NostrConnectEvent>(UNLIMITED)
val subId = newSubId()
@@ -25,6 +25,7 @@ import com.vitorpamplona.amethyst.commons.tor.TorType
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.isLocalHost
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.isOnion
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.TcpNoDelaySocketFactory
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.flow.SharingStarted
import kotlinx.coroutines.flow.StateFlow
@@ -66,6 +67,9 @@ class DesktopHttpClient(
private val directClient: OkHttpClient by lazy {
OkHttpClient
.Builder()
// TCP_NODELAY: keeps CLOSE-then-REQ bursts (feed switches) from
// nagling behind the unACKed CLOSE — see TcpNoDelaySocketFactory.
.socketFactory(TcpNoDelaySocketFactory)
.connectionPool(sharedConnectionPool)
.connectTimeout(BASE_TIMEOUT_SECONDS, TimeUnit.SECONDS)
.readTimeout(BASE_TIMEOUT_SECONDS, TimeUnit.SECONDS)
@@ -28,6 +28,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.relay.server.NostrServer
import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebsocketBuilder
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.BasicOkHttpWebSocket
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.TcpNoDelaySocketFactory
import com.vitorpamplona.quartz.nip01Core.store.IEventStore
import com.vitorpamplona.quartz.utils.Log
import com.vitorpamplona.quartz.utils.TimeUtils
@@ -117,6 +118,7 @@ class MirrorWorker(
if (websocketBuilder == null) {
OkHttpClient
.Builder()
.socketFactory(TcpNoDelaySocketFactory)
.pingInterval(Duration.ofSeconds(PING_INTERVAL_SECS))
.build()
} else {
@@ -0,0 +1,324 @@
/*
* Copyright (c) 2025 Vitor Pamplona
*
* Permission is hereby granted, free of charge, to any person obtaining a copy of
* this software and associated documentation files (the "Software"), to deal in
* the Software without restriction, including without limitation the rights to use,
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
* Software, and to permit persons to whom the Software is furnished to do so,
* subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all
* copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
*/
package com.vitorpamplona.geode.perf
import com.vitorpamplona.geode.KtorRelay
import com.vitorpamplona.geode.RelayEngine
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.TcpNoDelaySocketFactory
import com.vitorpamplona.quartz.nip01Core.store.sqlite.DefaultIndexingStrategy
import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
import com.vitorpamplona.quartz.utils.EventFactory
import io.ktor.server.application.install
import io.ktor.server.cio.CIO
import io.ktor.server.engine.embeddedServer
import io.ktor.server.routing.routing
import io.ktor.server.websocket.WebSockets
import io.ktor.server.websocket.webSocket
import io.ktor.websocket.Frame
import io.ktor.websocket.readText
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import okhttp3.OkHttpClient
import okhttp3.Request
import okhttp3.Response
import okhttp3.WebSocket
import okhttp3.WebSocketListener
import java.util.concurrent.ArrayBlockingQueue
import java.util.concurrent.TimeUnit
import kotlin.test.Test
import kotlin.test.assertTrue
/**
* Splits the wire-level small-REQ floor into transport vs relay work,
* over the production stack (Ktor CIO server, OkHttp client, loopback).
*
* - echo-1 / echo-22: a bare Ktor CIO websocket replying with 1 / 22
* ~250 B frames — the transport stack's round-trip floor; the burst
* variants show per-frame cost is negligible.
* - echo external / ext+1ms: replies produced from a foreign coroutine,
* immediately or after the connection loop parks — proves cross-context
* sends into CIO are prompt.
* - geode REQ / NOTICE / empty REQ / inproc: the relay's own layers.
*
* Historical note — this benchmark found the TcpNoDelaySocketFactory
* issue: without TCP_NODELAY on the client, a CLOSE (which relays never
* answer) followed by a REQ nagles the REQ behind the unACKed CLOSE for
* the ~40 ms delayed-ACK window, measured here as a flat 43.7 ms per
* round. With the factory (now used by every production client), geode's
* ~21-row REQ costs ~1.2 ms on the wire — matching relayBench — of which
* ~0.6 ms is the transport floor and ~0.5 ms the per-REQ server work
* documented in quartz/plans/2026-07-04-small-req-floor.md.
*/
class WireReqFloorBenchmark {
companion object {
const val EVENTS = 50_000
const val AUTHORS = 2_500
const val ROUNDS = 300
const val WARMUP = 50
const val BURST = 22
}
private fun hexId(seed: Int): String = seed.toString(16).padStart(64, '0')
private fun pubkey(seed: Int): String = (seed % AUTHORS).toString(16).padStart(64, 'a')
private val sig = "0".repeat(128)
private fun event(seed: Int): Event =
EventFactory.create(
id = hexId(seed),
pubKey = pubkey(seed),
createdAt = 1_600_000_000L + (seed * 7919) % 1_000_000,
kind = 1,
tags = emptyArray(),
content = "wire req floor benchmark $seed",
sig = sig,
)
private fun median(samples: LongArray): Double {
samples.sort()
return samples[samples.size / 2] / 1e6
}
/** Blocking one-connection driver: send [request], await [expectedFrames] replies. */
private class Driver(
client: OkHttpClient,
url: String,
) {
private val frames = ArrayBlockingQueue<String>(4096)
val socket: WebSocket
init {
val opened = ArrayBlockingQueue<Boolean>(1)
socket =
client.newWebSocket(
Request.Builder().url(url).build(),
object : WebSocketListener() {
override fun onOpen(
webSocket: WebSocket,
response: Response,
) {
opened.put(true)
}
override fun onMessage(
webSocket: WebSocket,
text: String,
) {
frames.put(text)
}
},
)
check(opened.poll(10, TimeUnit.SECONDS) == true) { "websocket did not open" }
}
var lastFirstFrameNanos: Long = 0
private set
fun roundTrip(
request: String,
expectedFrames: Int,
): Long {
val t0 = System.nanoTime()
socket.send(request)
repeat(expectedFrames) { i ->
checkNotNull(frames.poll(10, TimeUnit.SECONDS)) { "timed out waiting for frame" }
if (i == 0) lastFirstFrameNanos = System.nanoTime() - t0
}
return System.nanoTime() - t0
}
}
@Test
fun transportVsRelayFloor() =
runBlocking {
// TCP_NODELAY, via the same factory production clients use:
// without it, a CLOSE (never answered) followed by a REQ nagles
// the REQ behind the unACKed CLOSE for the ~40 ms delayed-ACK
// window — this benchmark measured a flat 43.7 ms per round
// before the factory existed, which is how the issue was found.
val client = OkHttpClient.Builder().socketFactory(TcpNoDelaySocketFactory).build()
// --- bare Ktor CIO echo server: N frames of ~250 B per request ---
val payload = "x".repeat(250)
val echo =
embeddedServer(CIO, port = 0) {
install(WebSockets)
routing {
webSocket("/") {
for (frame in incoming) {
if (frame !is Frame.Text) continue
val text = frame.readText()
if (text.startsWith("x")) {
// Reply from an EXTERNAL coroutine — the
// cross-context path geode's launched REQ
// handler uses.
val n = text.drop(1).toIntOrNull() ?: 1
launch(Dispatchers.Default) {
repeat(n) { outgoing.send(Frame.Text(payload)) }
}
} else if (text.startsWith("d")) {
// Same, but AFTER the connection loop has
// gone idle — does a parked CIO connection
// pick up a cross-thread send promptly?
val n = text.drop(1).toIntOrNull() ?: 1
launch(Dispatchers.Default) {
delay(1)
repeat(n) { outgoing.send(Frame.Text(payload)) }
}
} else {
val n = text.toIntOrNull() ?: 1
repeat(n) { outgoing.send(Frame.Text(payload)) }
}
}
}
}
}
echo.start(wait = false)
val echoPort =
echo.engine
.resolvedConnectors()
.first()
.port
val echoDriver = Driver(client, "ws://127.0.0.1:$echoPort/")
repeat(WARMUP) { echoDriver.roundTrip("1", 1) }
val echo1 = LongArray(ROUNDS) { echoDriver.roundTrip("1", 1) }
repeat(WARMUP) { echoDriver.roundTrip("$BURST", BURST) }
val echoN = LongArray(ROUNDS) { echoDriver.roundTrip("$BURST", BURST) }
repeat(WARMUP) { echoDriver.roundTrip("x1", 1) }
val echoExternal = LongArray(ROUNDS) { echoDriver.roundTrip("x1", 1) }
repeat(WARMUP) { echoDriver.roundTrip("d1", 1) }
val echoDelayed = LongArray(ROUNDS) { echoDriver.roundTrip("d1", 1) }
// --- geode: real REQ over the same stack ---
// Same store setup SmallReqFloorBenchmark proved answers this
// filter in ~0.18 ms in-process — this benchmark attributes
// the wire-vs-in-process delta, so the storage side must be
// the known-fast configuration.
val store = EventStore(dbName = null, indexStrategy = DefaultIndexingStrategy(indexEventsByPubkeyAlone = true))
(1..EVENTS).chunked(2000).forEach { chunk -> store.batchInsert(chunk.map { event(it) }) }
val relay = RelayEngine(url = "ws://127.0.0.1:7795/".normalizeRelayUrl(), store = store, parentContext = Dispatchers.IO + SupervisorJob())
val server = KtorRelay(relay, host = "127.0.0.1", port = 0).start()
val geodeDriver = Driver(client, server.url)
fun req(round: Int) = """["REQ","w$round",{"authors":["${pubkey(round)}"],"kinds":[1],"limit":50}]"""
// Each REQ answers with (rows + 1) frames — the EVENTs then
// EOSE — and the row count is author-dependent, so resolve
// every round's expected count from the store up front. Subs
// are CLOSEd after each round (silently, per NIP-01) so the
// relay's live-subscription registry doesn't grow with the
// round count and skew later samples.
suspend fun rowsFor(round: Int): Int = store.query<Event>(Filter(authors = listOf(pubkey(round)), kinds = listOf(1), limit = 50)).size
repeat(WARMUP) { round ->
val rows = rowsFor(round)
geodeDriver.roundTrip(req(round), rows + 1)
geodeDriver.socket.send("""["CLOSE","w$round"]""")
}
val geodeSamples = LongArray(ROUNDS)
val geodeFirst = LongArray(ROUNDS)
var totalRows = 0
for (i in 0 until ROUNDS) {
val round = WARMUP + i
val rows = rowsFor(round)
totalRows += rows
geodeSamples[i] = geodeDriver.roundTrip(req(round), rows + 1)
geodeFirst[i] = geodeDriver.lastFirstFrameNanos
geodeDriver.socket.send("""["CLOSE","w$round"]""")
}
// Probe A: malformed frame -> NOTICE, answered inline on the
// pump coroutine (no launch, no SQL). Probe B: REQ matching
// nothing -> lone EOSE (launch + SQL, no rows).
repeat(WARMUP) { geodeDriver.roundTrip("not json", 1) }
val notice = LongArray(ROUNDS) { geodeDriver.roundTrip("not json", 1) }
repeat(WARMUP) { i ->
geodeDriver.roundTrip("""["REQ","n$i",{"ids":["${"f".repeat(64)}"]}]""", 1)
geodeDriver.socket.send("""["CLOSE","n$i"]""")
}
val emptyReq =
LongArray(ROUNDS) { i ->
val t = geodeDriver.roundTrip("""["REQ","m$i",{"ids":["${"f".repeat(64)}"]}]""", 1)
geodeDriver.socket.send("""["CLOSE","m$i"]""")
t
}
// Same probe with NO CLOSE between rounds: if this is fast, the
// 44 ms is the client's CLOSE frame sitting unACKed (server
// sends nothing for a CLOSE) and Nagle holding the next REQ
// behind it until the ~40 ms delayed ACK — a client-side TCP
// artifact, not the relay.
repeat(WARMUP) { i -> geodeDriver.roundTrip("""["REQ","nc$i",{"ids":["${"a".repeat(64)}"]}]""", 1) }
val emptyNoClose =
LongArray(ROUNDS) { i ->
geodeDriver.roundTrip("""["REQ","ncm$i",{"ids":["${"a".repeat(64)}"]}]""", 1)
}
// Probe C: same RelayEngine, no wire — an in-process session
// while Ktor keeps serving. Splits engine-state issues from
// transport-adjacent ones.
val inproc = LongArray(ROUNDS)
run {
val q = ArrayBlockingQueue<String>(4096)
val session = relay.server.connect { q.put(it) }
repeat(WARMUP) { i ->
session.receive("""["REQ","p$i",{"ids":["${"e".repeat(64)}"]}]""")
checkNotNull(q.poll(10, TimeUnit.SECONDS))
session.receive("""["CLOSE","p$i"]""")
}
for (i in 0 until ROUNDS) {
val t0 = System.nanoTime()
session.receive("""["REQ","q$i",{"ids":["${"e".repeat(64)}"]}]""")
checkNotNull(q.poll(10, TimeUnit.SECONDS))
inproc[i] = System.nanoTime() - t0
session.receive("""["CLOSE","q$i"]""")
}
session.close()
}
assertTrue(totalRows > 0)
println("WireReqFloorBenchmark (loopback, OkHttp client) @ ${EVENTS / 1000}k events, medians of $ROUNDS")
println(" echo 1 frame: ${"%6.3f".format(median(echo1))} ms")
println(" echo $BURST frames: ${"%6.3f".format(median(echoN))} ms")
println(" echo 1 frame (external): ${"%6.3f".format(median(echoExternal))} ms")
println(" echo 1 frame (ext+1ms): ${"%6.3f".format(median(echoDelayed))} ms")
println(" geode REQ (~${totalRows / ROUNDS} rows+EOSE): ${"%6.3f".format(median(geodeSamples))} ms (first frame ${"%6.3f".format(median(geodeFirst))} ms)")
println(" geode NOTICE (inline): ${"%6.3f".format(median(notice))} ms")
println(" geode empty REQ (launch): ${"%6.3f".format(median(emptyReq))} ms")
println(" geode empty REQ no CLOSE: ${"%6.3f".format(median(emptyNoClose))} ms")
println(" geode empty REQ (inproc): ${"%6.3f".format(median(inproc))} ms")
geodeDriver.socket.close(1000, null)
echoDriver.socket.close(1000, null)
server.stop(0, 1_000)
relay.close()
echo.stop(0, 500)
client.dispatcher.executorService.shutdown()
}
}
+31 -12
View File
@@ -42,17 +42,36 @@ the container drift band** — strfry's own numbers drifted ±30% run to
run, and inline-eligible scenarios moved the same as ineligible ones.
Reverted per the keep-only-winners rule.
## Where the floor actually is
## Where the floor actually is (corrected after WireReqFloorBenchmark)
In-process REQ→EOSE is ~0.5–0.6 ms, but the wire-level p50 is
1.2–1.7 ms: the missing ~1 ms per REQ is transport-side — the Ktor CIO
frame write path, per-frame sends with no batching, and the client
round trip — which is backlog item 6 (websocket send path, 13–25% of
ingest CPU in the JFR profile) territory. strfry completes the whole
round trip in under 0.6 ms on a single event loop with uWebSockets.
The follow-up wire benchmark (geode's `WireReqFloorBenchmark`, Ktor CIO
+ OkHttp on loopback) attributed the full path:
**Do not retry** coroutine-dispatch shaving for this gap without first
measuring the transport side: instrument the time between
`RelaySession`'s `onSend` invocation and the frame actually leaving
the socket, and compare permessage-deflate/frame-batching settings
against strfry's uWS configuration.
| leg | ms |
|---|---:|
| bare Ktor CIO echo round trip (1 or 22 frames — same) | ~0.9–1.0 |
| geode NOTICE (inline, full pump + Ktor send) | ~0.5 |
| geode empty REQ (launch + SQL, 0 rows) | ~0.6–0.8 |
| geode ~21-row REQ, wire | ~1.25 (= relayBench's number) |
**geode's websocket send path has no latency problem** — per-frame burst
cost is negligible (echo-22 ≈ echo-1), the pump adds ~nothing (NOTICE ≈
0.5 ms), and the residual vs strfry (~0.5 ms/REQ) is the per-REQ server
work already investigated above. Frame batching / permessage-deflate
would not move these numbers. Backlog item 6's remaining open angle is
the INGEST-side CPU share (13–25% in the JFR profile) — a throughput
question, not this latency one.
**The real find was client-side.** The first wire measurements showed a
flat 43.7 ms per REQ — which turned out to be the benchmark's own OkHttp
client: OkHttp does not set TCP_NODELAY, and the CLOSE-then-REQ pattern
(every feed/filter switch!) nagles the REQ behind the unACKed CLOSE
(relays never answer CLOSE) for the ~40 ms delayed-ACK window.
relayBench's harness client already carried a no-delay socket factory —
which is why bench numbers never showed it — but the production clients
(Android relay pool, Desktop, amy, geode's mirror) did not. Fixed by
`TcpNoDelaySocketFactory` (quartz jvmAndroid), now used by all of them.
**Do not retry** relay-side latency work for the small-REQ gap; the
addressable remainder is the ~0.4 ms of per-REQ dispatch machinery this
doc's revert already covers, and it does not show on the wire.
@@ -0,0 +1,83 @@
/*
* Copyright (c) 2025 Vitor Pamplona
*
* Permission is hereby granted, free of charge, to any person obtaining a copy of
* this software and associated documentation files (the "Software"), to deal in
* the Software without restriction, including without limitation the rights to use,
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
* Software, and to permit persons to whom the Software is furnished to do so,
* subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all
* copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
*/
package com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp
import java.net.InetAddress
import java.net.InetSocketAddress
import java.net.Socket
import javax.net.SocketFactory
/**
* [SocketFactory] that enables `TCP_NODELAY` — pass to
* `OkHttpClient.Builder.socketFactory` for every client that holds Nostr
* relay WebSockets.
*
* OkHttp does not disable Nagle's algorithm, and a Nostr client's hottest
* wire pattern trips over it: CLOSE immediately followed by REQ — every
* feed or filter switch. Relays never reply to a CLOSE (NIP-01), so its
* bytes sit unACKed for the peer's delayed-ACK window, and Nagle holds
* the REQ back behind them until that ACK arrives. Measured on loopback
* as a flat ~44 ms added to every such REQ's round trip (geode's
* `WireReqFloorBenchmark`, which is how this was found); WAN delayed-ACK
* windows are the same order. relayBench's harness client ships the same
* fix, which is why benchmark numbers never showed it.
*
* OkHttp only uses this factory for direct connections — proxied (e.g.
* Tor SOCKS) connections take a different path and keep their transport's
* defaults.
*/
object TcpNoDelaySocketFactory : SocketFactory() {
private fun socket() = Socket().apply { tcpNoDelay = true }
override fun createSocket(): Socket = socket()
override fun createSocket(
host: String?,
port: Int,
): Socket = socket().apply { connect(InetSocketAddress(host, port)) }
override fun createSocket(
host: String?,
port: Int,
localHost: InetAddress?,
localPort: Int,
): Socket =
socket().apply {
bind(InetSocketAddress(localHost, localPort))
connect(InetSocketAddress(host, port))
}
override fun createSocket(
host: InetAddress?,
port: Int,
): Socket = socket().apply { connect(InetSocketAddress(host, port)) }
override fun createSocket(
address: InetAddress?,
port: Int,
localAddress: InetAddress?,
localPort: Int,
): Socket =
socket().apply {
bind(InetSocketAddress(localAddress, localPort))
connect(InetSocketAddress(address, port))
}
}