feat: add negentropySyncFanOut — multi-connection sync from one relay

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 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018saXqYfAa3RvSJoDXK591R
This commit is contained in:
Claude
2026-07-03 04:27:46 +00:00
parent 5a993810bb
commit 4bbfd86a56
6 changed files with 487 additions and 5 deletions
@@ -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
@@ -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,
)
@@ -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<INostrClient>,
relay: NormalizedRelayUrl,
filter: Filter,
localEntries: List<IdAndTime> = 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<List<HexKey>>(clients.size * reqsPerClient * 2)
// Download → consumer funnel; single consumer keeps onEvent
// single-threaded and the maxEvents cap exact.
val events = Channel<Event>(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
@@ -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<INostrClient>,
relay: NormalizedRelayUrl,
filter: Filter,
localEntries: List<IdAndTime>,
@@ -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<HexKey>,
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
@@ -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<NostrClient>()
private fun clients(count: Int): List<NostrClient> =
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<Event> = 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<Event>()
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<Event>()
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<Event>()
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<Event>()
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)
}
}
@@ -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