From fd6662ca30fb5e8c9c9c263267e5640c5577a41c Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 4 Jul 2026 00:38:09 +0000 Subject: [PATCH] =?UTF-8?q?perf(quartz):=20TCP=5FNODELAY=20for=20every=20r?= =?UTF-8?q?elay=20websocket=20client=20=E2=80=94=20kills=20CLOSE=E2=86=92R?= =?UTF-8?q?EQ=20Nagle=20stalls?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 Claude-Session: https://claude.ai/code/session_01TtDNpayEYvJH7QuPswND3A --- .../okhttp/OkHttpClientFactoryForRelays.kt | 7 + .../com/vitorpamplona/amethyst/cli/Context.kt | 3 +- .../amethyst/cli/commands/NipCommand.kt | 3 +- .../amethyst/cli/commands/NostrConnect.kt | 3 +- .../desktop/network/DesktopHttpClient.kt | 4 + .../geode/mirror/MirrorWorker.kt | 2 + .../geode/perf/WireReqFloorBenchmark.kt | 324 ++++++++++++++++++ quartz/plans/2026-07-04-small-req-floor.md | 43 ++- .../sockets/okhttp/TcpNoDelaySocketFactory.kt | 83 +++++ 9 files changed, 457 insertions(+), 15 deletions(-) create mode 100644 geode/src/test/kotlin/com/vitorpamplona/geode/perf/WireReqFloorBenchmark.kt create mode 100644 quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/TcpNoDelaySocketFactory.kt diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpClientFactoryForRelays.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpClientFactoryForRelays.kt index 7704771de1..8b27d5eee3 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpClientFactoryForRelays.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/okhttp/OkHttpClientFactoryForRelays.kt @@ -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) diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt index a96fc68f69..4ffe1a77fc 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt @@ -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( diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/NipCommand.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/NipCommand.kt index 665c4a039c..446c53fd48 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/NipCommand.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/NipCommand.kt @@ -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> { 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() diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/NostrConnect.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/NostrConnect.kt index c26c166854..08ffe57337 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/NostrConnect.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/NostrConnect.kt @@ -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(UNLIMITED) val subId = newSubId() diff --git a/desktopApp/src/jvmMain/kotlin/com/vitorpamplona/amethyst/desktop/network/DesktopHttpClient.kt b/desktopApp/src/jvmMain/kotlin/com/vitorpamplona/amethyst/desktop/network/DesktopHttpClient.kt index a63091f7a9..6f4c290668 100644 --- a/desktopApp/src/jvmMain/kotlin/com/vitorpamplona/amethyst/desktop/network/DesktopHttpClient.kt +++ b/desktopApp/src/jvmMain/kotlin/com/vitorpamplona/amethyst/desktop/network/DesktopHttpClient.kt @@ -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) diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt index d12fffc496..f87737d89c 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt @@ -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 { diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/perf/WireReqFloorBenchmark.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/perf/WireReqFloorBenchmark.kt new file mode 100644 index 0000000000..b0778a2d07 --- /dev/null +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/perf/WireReqFloorBenchmark.kt @@ -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(4096) + val socket: WebSocket + + init { + val opened = ArrayBlockingQueue(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(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(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() + } +} diff --git a/quartz/plans/2026-07-04-small-req-floor.md b/quartz/plans/2026-07-04-small-req-floor.md index 7723eba0c5..77ea84fe55 100644 --- a/quartz/plans/2026-07-04-small-req-floor.md +++ b/quartz/plans/2026-07-04-small-req-floor.md @@ -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. diff --git a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/TcpNoDelaySocketFactory.kt b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/TcpNoDelaySocketFactory.kt new file mode 100644 index 0000000000..2f175a3b3b --- /dev/null +++ b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip01Core/relay/sockets/okhttp/TcpNoDelaySocketFactory.kt @@ -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)) + } +}