From 4bbfd86a5689b2750e8242125273ab415ec8397b Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 3 Jul 2026 04:27:46 +0000 Subject: [PATCH] =?UTF-8?q?feat:=20add=20negentropySyncFanOut=20=E2=80=94?= =?UTF-8?q?=20multi-connection=20sync=20from=20one=20relay?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit One reconcile feeds by-id download batches to N clients (one socket each) x reqsPerClient workers; reconcile windows also round-robin across the connections, since a single connection produced need-ids at only ~9k/s on a 2.6M corpus and starved the downloads. Events funnel through a bounded channel to a single consumer (exact maxEvents, single-threaded onEvent); all stages backpressure; localEntries diffing and have-counting match negentropyReconcile. reconcileWindows became multi-client internally; single-client paths pass listOf(this). Production shootout (same-run pairs, 100k cap): +18% to +64% over the tuned single client, capped by the relay's server-side reconcile id production rather than download parallelism (which the by-id matrix shows scales 2.7x with connections). 4 in-process multi-client tests; 18 negentropy tests green. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_018saXqYfAa3RvSJoDXK591R --- .../2026-07-02-nostrclient-receiver-perf.md | 27 +++ .../accessories/NegentropyFanOutResult.kt | 39 ++++ .../NostrClientNegentropyFanOutExt.kt | 198 ++++++++++++++++++ .../NostrClientNegentropySyncExt.kt | 17 +- .../relay/NostrClientNegentropyFanOutTest.kt | 147 +++++++++++++ .../relay/prodbench/BulkDownloadBenchmark.kt | 64 ++++++ 6 files changed, 487 insertions(+), 5 deletions(-) create mode 100644 quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropyFanOutResult.kt create mode 100644 quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropyFanOutExt.kt create mode 100644 quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/NostrClientNegentropyFanOutTest.kt diff --git a/quartz/plans/2026-07-02-nostrclient-receiver-perf.md b/quartz/plans/2026-07-02-nostrclient-receiver-perf.md index d78dcc3433..1cd0f676ae 100644 --- a/quartz/plans/2026-07-02-nostrclient-receiver-perf.md +++ b/quartz/plans/2026-07-02-nostrclient-receiver-perf.md @@ -535,6 +535,33 @@ submission-order preservation, preVerified short-circuit, callback-crash resilience. Wiring it into `CacheClientConnector`/`LocalCache` (consume with `wasVerified = true`) is the app-side follow-up. +### negentropySyncFanOut: multi-connection sync (done) + +`accessories/NostrClientNegentropyFanOutExt`: one reconcile feeding by-id +download batches to N clients (one socket each) × `reqsPerClient` workers, +with reconcile windows ALSO round-robined across the connections. Events +funnel through a bounded channel to a single consumer (exact `maxEvents`, +single-threaded `onEvent`); everything backpressures; `localEntries` diffing +and have-reporting work as in `negentropyReconcile`. 4 in-process +multi-client tests (full download, cap, local-diff, single-client +degradation) — 18 negentropy tests green. + +Production shootout (100k cap, kind 30382, same-run pairs): fan-out 4×8 vs +tuned single client measured **+18% to +64%** (7.3k vs 6.1k/s; 7.5k vs +4.5k/s) with heavy relay-side variance — the relay is being actively +backfilled and its NEG snapshot caching favors whichever variant runs +second. The honest bound: **this relay produces reconcile ids at only +~9–10k/s regardless of connection count** (server-side snapshot build), so +end-to-end fan-out gains cap there even though the download stage alone +scales 2.7× with connections (the by-id matrix). On download-bound syncs +(bigger events, relays with faster NEG production, phone-CPU clients) the +gain approaches the matrix ratio. + +**Descoped with evidence — adaptive in-flight control (#6):** an AIMD window +on download REQs would idle-down against today's reconcile-bound relay and +the benchmarks could not demonstrate a win (the criterion for shipping); +revisit if a relay shows download-bound behavior with volatile capacity. + ## Recommendations (in order of value/risk) 1. **Move Schnorr verification off the receiver coroutine** in the app's diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropyFanOutResult.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropyFanOutResult.kt new file mode 100644 index 0000000000..96392ace8a --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NegentropyFanOutResult.kt @@ -0,0 +1,39 @@ +/* + * 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.client.accessories + +/** + * Outcome of a [negentropySyncFanOut] run. + * + * @property needCount ids the relay reported missing locally (full set size when uncapped). + * @property haveCount ids the local set has that the relay lacks (only counted when + * [negentropySyncFanOut]'s `localEntries` is non-empty). + * @property downloaded distinct events delivered through `onEvent`. + * @property windows `created_at` windows the reconcile split into. + * @property connections how many clients actually served download batches. + */ +class NegentropyFanOutResult( + val needCount: Int, + val haveCount: Int, + val downloaded: Int, + val windows: Int, + val connections: Int, +) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropyFanOutExt.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropyFanOutExt.kt new file mode 100644 index 0000000000..cce9d24bf1 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropyFanOutExt.kt @@ -0,0 +1,198 @@ +/* + * 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.client.accessories + +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.core.HexKey +import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient +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.store.IdAndTime +import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.coroutineScope +import kotlinx.coroutines.ensureActive +import kotlinx.coroutines.joinAll +import kotlinx.coroutines.launch +import kotlin.concurrent.atomics.AtomicInt +import kotlin.concurrent.atomics.ExperimentalAtomicApi +import kotlin.concurrent.atomics.incrementAndFetch +import kotlin.coroutines.coroutineContext + +/** + * [negentropySync] scaled past the single-connection ceiling: ONE reconcile + * (on `clients[0]`) feeding by-id download batches to EVERY client in + * [clients] — each client being its own socket to the same [relay]. + * + * Why this shape: a relay serves a single connection only so fast — the + * production by-id matrix measured ~5.5k events/s per connection (saturating + * at ~8 in-flight REQs) but ~15k events/s across 4 connections; see + * `quartz/plans/2026-07-02-nostrclient-receiver-perf.md`. Reconciliation is + * cheap and stays on one connection; the downloads are what need the fan-out. + * + * Mechanics: + * - `clients[0]` runs [negentropyReconcile] (windows, overflow splits, + * optional [localEntries] diffing — all identical semantics), streaming + * need-id batches into a bounded queue; + * - every client runs [reqsPerClient] download workers pulling batches off + * that shared queue ([fetchByIds]: one REQ per batch, collected to EOSE); + * - events funnel through a bounded channel to a single consumer, so + * [onEvent] runs single-threaded and [maxEvents] is exact; the bounded + * stages back-pressure all the way to the reconcile. + * + * Reconcile windows round-robin across the clients too (a single connection + * produced need-ids at only ~9k/s on a 2.6M corpus and starved the download + * workers). Subscription budget per connection is the caller's job (relays + * cap concurrent subs; strfry defaults to 20): each client holds at peak + * `reqsPerClient + ceil(reconcileConcurrency / clients.size) + 1`. + * + * Passing a single-element [clients] degrades gracefully to (pipelined) + * [negentropySync] behavior. Duplicate clients in the list would just share + * a connection — use distinct instances. + * + * @throws NegentropySyncException when a window cannot be reconciled (same + * contract as [negentropySync]); the whole fan-out is cancelled. + */ +@OptIn(ExperimentalAtomicApi::class) +suspend fun negentropySyncFanOut( + clients: List, + relay: NormalizedRelayUrl, + filter: Filter, + localEntries: List = emptyList(), + maxEvents: Int = 0, + reqsPerClient: Int = 10, + fetchBatch: Int = 250, + idleTimeoutMs: Long = 120_000L, + reconcileConcurrency: Int = 2, + onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, + onEvent: (Event) -> Unit, +): NegentropyFanOutResult { + require(clients.isNotEmpty()) { "at least one client is required" } + + val need = AtomicInt(0) + val have = AtomicInt(0) + val windows = AtomicInt(0) + val used = AtomicInt(0) + var downloaded = 0 + + // Pin every client's connection open for the whole run: between two + // batches (or two reconcile windows) a client briefly has no live + // subscription and its pool would otherwise drop the relay. + val keepAlives = + clients.map { client -> + val subId = newSubId() + client.subscribe(subId, mapOf(relay to listOf(Filter(ids = listOf(KEEP_ALIVE_ID)))), null) + client to subId + } + + try { + coroutineScope { + // Reconcile → download handoff. Bounded so download saturation + // back-pressures the reconcile instead of piling up ids. + val idBatches = Channel>(clients.size * reqsPerClient * 2) + + // Download → consumer funnel; single consumer keeps onEvent + // single-threaded and the maxEvents cap exact. + val events = Channel(DELIVERY_BUFFER) + + val workers = + clients.flatMap { client -> + List(reqsPerClient.coerceAtLeast(1)) { + launch { + var servedAny = false + for (batch in idBatches) { + coroutineContext.ensureActive() + if (!servedAny) { + servedAny = true + used.incrementAndFetch() + } + for (event in client.fetchByIds(relay, batch, idleTimeoutMs)) { + events.send(event) + } + } + } + } + } + + val producer = + launch { + try { + val sorted = + if (localEntries.size > 1) localEntries.sortedBy { it.createdAt } else localEntries + // Windows reconcile round-robin ACROSS the clients so + // server-side snapshot builds parallelize per connection + // (a single connection produced ids at only ~9k/s and + // starved the download workers). + reconcileWindows( + clients = clients, + relay = relay, + filter = filter, + localEntries = sorted, + idleTimeoutMs = idleTimeoutMs, + batchSize = fetchBatch, + reconcileConcurrency = reconcileConcurrency, + onWindow = { windows.incrementAndFetch() }, + onNeed = { + need.addAndFetch(it) + }, + onHave = { have.addAndFetch(it) }, + sendNeedBatch = { batch -> idBatches.send(batch) }, + sendHaveBatch = if (localEntries.isEmpty()) null else { _ -> }, + ) + } finally { + idBatches.close() + } + } + + val closer = + launch { + producer.join() + workers.joinAll() + events.close() + } + + for (event in events) { + downloaded++ + onEvent(event) + onProgress?.invoke(need.load(), downloaded) + if (maxEvents in 1..downloaded) break + } + + // Cap reached (or producer done): stop everything still running. + producer.cancel() + workers.forEach { it.cancel() } + closer.cancel() + } + } finally { + keepAlives.forEach { (client, subId) -> client.unsubscribe(subId) } + } + + return NegentropyFanOutResult( + needCount = need.load(), + haveCount = have.load(), + downloaded = downloaded, + windows = windows.load(), + connections = used.load(), + ) +} + +/** Bounded buffer between download workers and the single delivery consumer. */ +private const val DELIVERY_BUFFER = 256 diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt index 9e846a91fd..9d3a65eeee 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt @@ -410,6 +410,7 @@ private suspend fun INostrClient.syncPipeline( } reconcileWindows( + clients = listOf(this@syncPipeline), relay = relay, filter = filter, localEntries = emptyList(), @@ -441,7 +442,11 @@ private suspend fun INostrClient.syncPipeline( * failure cancels the whole scope. */ @OptIn(ExperimentalAtomicApi::class) -private suspend fun INostrClient.reconcileWindows( +internal suspend fun reconcileWindows( + // one or more connections to the SAME relay; reconciler i runs its NEG + // sessions on clients[i % size], so concurrent windows spread across + // connections (server-side snapshot builds are paced per connection) + clients: List, relay: NormalizedRelayUrl, filter: Filter, localEntries: List, @@ -465,13 +470,14 @@ private suspend fun INostrClient.reconcileWindows( pending.send(filter) val reconcilers = - List(reconcileConcurrency.coerceAtLeast(1)) { + List(reconcileConcurrency.coerceAtLeast(1)) { reconcilerIndex -> launch { + val client = clients[reconcilerIndex % clients.size] for (window in pending) { coroutineContext.ensureActive() val outcome = - reconcileStreaming( + client.reconcileStreaming( relay = relay, filter = window, localEntries = entriesForWindow(localEntries, window.since, window.until), @@ -650,6 +656,7 @@ suspend fun INostrClient.negentropyReconcile( } reconcileWindows( + clients = listOf(this), relay = relay, filter = filter, localEntries = sorted, @@ -934,7 +941,7 @@ private fun isOverflow(reason: String): Boolean = * replay the batch; without this the same event would be delivered twice. We rely on * NIP-77 yielding a distinct id set across batches, so no global dedup is needed. */ -private suspend fun INostrClient.fetchByIds( +internal suspend fun INostrClient.fetchByIds( relay: NormalizedRelayUrl, batch: List, idleTimeoutMs: Long, @@ -1009,7 +1016,7 @@ private const val DELIVERY_BUFFER = 256 * duration of a sync. Synthetic/real event ids are SHA-256 digests, so this never * collides with an actual event. */ -private val KEEP_ALIVE_ID = "f".repeat(64) +internal val KEEP_ALIVE_ID = "f".repeat(64) /** * Finite fallback bounds (ms) for the two waits that must stay bounded even when the diff --git a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/NostrClientNegentropyFanOutTest.kt b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/NostrClientNegentropyFanOutTest.kt new file mode 100644 index 0000000000..371da758d3 --- /dev/null +++ b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/NostrClientNegentropyFanOutTest.kt @@ -0,0 +1,147 @@ +/* + * 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 + +import com.vitorpamplona.geode.fixtures.SyntheticEvents +import com.vitorpamplona.geode.testing.RelayClientTest +import com.vitorpamplona.geode.testing.preload +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient +import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.negentropySyncFanOut +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.store.IdAndTime +import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.withTimeout +import org.junit.After +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertTrue + +class NostrClientNegentropyFanOutTest : RelayClientTest() { + // extra connections to the same in-process relay + private val extraClients = mutableListOf() + + private fun clients(count: Int): List = + buildList { + add(client) + repeat(count - 1) { + add(NostrClient(hub, scope).also { extraClients.add(it) }) + } + } + + @After + fun tearDownExtraClients() { + extraClients.forEach { it.close() } + } + + private fun events(range: IntRange): List = range.map { SyntheticEvents.fakeEvent(idSeed = it, kind = 1, createdAt = it.toLong()) } + + @Test + fun fullDownloadAcrossMultipleConnections() = + runBlocking { + val seeded = events(1..60) + defaultRelay.preload(seeded) + + val got = mutableListOf() + val result = + withTimeout(30_000) { + negentropySyncFanOut( + clients = clients(3), + relay = defaultRelayUrl, + filter = Filter(kinds = listOf(1)), + fetchBatch = 5, + reqsPerClient = 2, + ) { got.add(it) } + } + + assertEquals(60, got.size, "every event delivered") + assertEquals(60, got.map { it.id }.toSet().size, "no duplicates") + assertEquals(60, result.needCount) + assertEquals(60, result.downloaded) + assertTrue(result.connections >= 1, "at least one connection served batches") + } + + @Test + fun maxEventsCapsDeliveryAndStopsTheFanOut() = + runBlocking { + defaultRelay.preload(events(1..50)) + + val got = mutableListOf() + val result = + withTimeout(30_000) { + negentropySyncFanOut( + clients = clients(2), + relay = defaultRelayUrl, + filter = Filter(kinds = listOf(1)), + maxEvents = 12, + fetchBatch = 4, + ) { got.add(it) } + } + + assertEquals(12, got.size) + assertEquals(12, result.downloaded) + } + + @Test + fun localEntriesSkipWhatWeAlreadyHaveAndReportHaves() = + runBlocking { + val relayOnly = events(1..20) + val shared = events(101..110) + val localOnly = events(201..205) + defaultRelay.preload(relayOnly + shared) + + val got = mutableListOf() + val result = + withTimeout(30_000) { + negentropySyncFanOut( + clients = clients(2), + relay = defaultRelayUrl, + filter = Filter(kinds = listOf(1)), + localEntries = (shared + localOnly).map { IdAndTime(it.createdAt, it.id) }, + fetchBatch = 6, + ) { got.add(it) } + } + + assertEquals(relayOnly.map { it.id }.toSet(), got.map { it.id }.toSet(), "only relay-only events download") + assertEquals(20, result.needCount) + assertEquals(5, result.haveCount, "local-only events reported as haves") + } + + @Test + fun singleClientDegradesGracefully() = + runBlocking { + defaultRelay.preload(events(1..25)) + + val got = mutableListOf() + val result = + withTimeout(30_000) { + negentropySyncFanOut( + clients = clients(1), + relay = defaultRelayUrl, + filter = Filter(kinds = listOf(1)), + fetchBatch = 10, + ) { got.add(it) } + } + + assertEquals(25, result.downloaded) + assertEquals(25, got.map { it.id }.toSet().size) + } +} diff --git a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/prodbench/BulkDownloadBenchmark.kt b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/prodbench/BulkDownloadBenchmark.kt index d4f6725dd2..387af8d37f 100644 --- a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/prodbench/BulkDownloadBenchmark.kt +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/prodbench/BulkDownloadBenchmark.kt @@ -28,6 +28,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.NegentropySyncException import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAllPages import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.negentropySync +import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.negentropySyncFanOut import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl @@ -482,6 +483,69 @@ class BulkDownloadBenchmark { httpClient.dispatcher.executorService.shutdown() } + @Test + fun negentropyFanOutShootout() { + if (System.getenv("PROD_RELAY_BENCH") == null && System.getProperty("prodRelayBench") == null) { + println("negentropyFanOutShootout skipped. Run with -PprodRelayBench=1 to enable.") + return + } + + val httpClient = + OkHttpClient + .Builder() + .connectTimeout(15, TimeUnit.SECONDS) + .readTimeout(120, TimeUnit.SECONDS) + .pingInterval(30, TimeUnit.SECONDS) + .build() + + val relay = PROD_RELAY.normalizeRelayUrl() + println("=== NEGENTROPY FAN-OUT SHOOTOUT: $PROD_RELAY kinds=[$PROD_KIND], cap 100k events ===") + + runBlocking { + // Baseline: one connection, tuned pipeline (measured ~4.5k/s). + negentropyVariant( + "single client (12r/4rc)", + httpClient, + maxEvents = 100_000, + maxConcurrentReqs = 12, + reconcileConcurrency = 4, + idBufferBatches = 96, + budgetMs = 240_000L, + ) + + // Fan-out: 4 connections, 10 download REQs each + reconcile on the first. + val clients = List(4) { NostrClient(BasicOkHttpWebSocket.Builder { httpClient }) } + try { + val count = AtomicLong(0) + val start = System.nanoTime() + val result = + withTimeoutOrNull(240_000L) { + negentropySyncFanOut( + clients = clients, + relay = relay, + filter = Filter(kinds = listOf(PROD_KIND)), + maxEvents = 100_000, + reqsPerClient = 8, + fetchBatch = 250, + reconcileConcurrency = 6, + ) { count.incrementAndGet() } + } + val wall = System.nanoTime() - start + report("fan-out 4 clients x 8 reqs", count.get(), 0, wall) + if (result == null) { + println(" (budget exceeded — partial)") + } else { + println(" windows=${result.windows} need=${result.needCount} connectionsUsed=${result.connections}") + } + } catch (e: Exception) { + println(" !! fan-out errored: ${e::class.simpleName} ${e.message}") + } finally { + clients.forEach { it.close() } + } + } + httpClient.dispatcher.executorService.shutdown() + } + // ------------------------------------------------------------------ // Per-connection wall: is a single connection limited by the relay's // pacing or by the client's serial parse? A raw socket that only scans