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)) + } +}