From 2fb44d81663a6f7d47cb6e05b9e115786ee154c9 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 3 Jul 2026 02:28:32 +0000 Subject: [PATCH] =?UTF-8?q?feat:=20pipeline=20negentropySync=20=E2=80=94?= =?UTF-8?q?=20global=20worker=20pool,=20concurrent=20window=20reconciles?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Restructures the sync around the measured bottlenecks: one download worker pool now spans the whole sync (windows no longer join before the next reconcile starts), overflow-split windows are reconciled by a caller-set number of concurrent NEG sessions (reconcileConcurrency, default 1) from a shared work queue, and the reconcile-to-download buffer depth is exposed (idBufferBatches). No NIP-11 auto-detection: peak subscription usage (maxConcurrentReqs + reconcileConcurrency + 1) is documented and budgeting it against the relay's max_subscriptions is the caller's call. All 44 negentropy tests pass unchanged. Production shootout on a 2.6M corpus, single connection: 3.6k events/s with old-equivalent params, 4.5k/s tuned (12 reqs + 4 reconcilers) — ~82% of the measured ~5.5k/s per-connection by-id ceiling, vs 1.6k/s for the old implementation. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_018saXqYfAa3RvSJoDXK591R --- .../2026-07-02-nostrclient-receiver-perf.md | 35 +++ .../NostrClientNegentropySyncExt.kt | 228 +++++++++++------- .../relay/prodbench/BulkDownloadBenchmark.kt | 91 +++++++ 3 files changed, 266 insertions(+), 88 deletions(-) diff --git a/quartz/plans/2026-07-02-nostrclient-receiver-perf.md b/quartz/plans/2026-07-02-nostrclient-receiver-perf.md index 2003a1fc9e..78321be8fc 100644 --- a/quartz/plans/2026-07-02-nostrclient-receiver-perf.md +++ b/quartz/plans/2026-07-02-nostrclient-receiver-perf.md @@ -373,6 +373,41 @@ in-flight until events/s stops improving / per-REQ latency inflates, then back off), since the same relay's sweet spot moved 8× in one day as its dataset grew. At today's ~15k/s ceiling: 2.6M events ≈ 3 min; 10M ≈ 11 min. +### negentropySync pipelining (implemented, 2026-07-03) + +`NostrClientNegentropySyncExt` was restructured around the findings — no +NIP-11 auto-detection, the budget knobs are the caller's: + +- **Global download-worker pool** spanning the whole sync (was per-window + with a join between windows, so the connection idled while the next + window's NEG rounds ran). +- **`reconcileConcurrency` parameter** — overflow-split windows are + reconciled by N concurrent NEG sessions from a shared work queue (was a + strictly sequential lo-half/hi-half recursion). Default 1. +- **`idBufferBatches` parameter** — depth of the bounded buffer between + reconciliation and download workers (was hardcoded to `workerCount`). +- KDoc contract: peak subscriptions = `maxConcurrentReqs + + reconcileConcurrency + 1`; sizing it under the relay's + `limitation.max_subscriptions` is the caller's responsibility. + +All 44 negentropy tests (unit + end-to-end against the in-process relay + +geode interop) pass unchanged. Production shootout on the 2.6M corpus, +100k-event cap, same day, single connection +(`BulkDownloadBenchmark.negentropyPipelineShootout`): + +| variant | events/s | note | +|---|---|---| +| old implementation (previous day, 31k corpus) | 1,575 | not same-day comparable | +| new pipeline, old-equivalent params (8 reqs, seq windows) | 3,601 | global pool already overlaps windows | +| new pipeline, tuned (12 reqs, 4 reconcilers, 96-batch buffer) | **4,497** | 17 subs ≤ strfry's 20 cap | + +Tuned reaches **~82% of the measured ~5.5k/s single-connection by-id +ceiling**; the 4 reconcilers also produced ids 2.3× faster (need=171,812 +enumerated vs 109,336 in less wall time). The remaining gap is reconcile +burstiness plus by-id query cost. Going past ~5.5k/s on this relay requires +multiple connections (caller-level: disjoint `created_at` windows per +client), which the matrix caps at ~15k/s relay-wide. + ### Answer for "10M events from one relay, fastest" At nosfabrica's measured page cadence, one connection ≈ 3.7k events/s → 10M 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 a5b218cf76..28f7b88481 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 @@ -43,6 +43,10 @@ import kotlinx.coroutines.joinAll import kotlinx.coroutines.launch import kotlinx.coroutines.withTimeoutOrNull import kotlin.concurrent.Volatile +import kotlin.concurrent.atomics.AtomicInt +import kotlin.concurrent.atomics.ExperimentalAtomicApi +import kotlin.concurrent.atomics.decrementAndFetch +import kotlin.concurrent.atomics.incrementAndFetch import kotlin.coroutines.coroutineContext import kotlin.math.min import kotlin.time.TimeSource @@ -121,9 +125,26 @@ class NegentropySyncResult( * it and the disconnect is turned into a clean abort. Pass `0` to disable the * watchdog entirely and run until the socket drops (download batches keep a finite * internal idle bound regardless, so a single stuck batch can't hang the pipeline). + * @param reconcileConcurrency how many `created_at` windows are reconciled at once + * after an over-cap split (each holds one NEG session on the connection). `1` + * reproduces the old strictly-sequential window walk. Raising it overlaps the + * reconcile round-trips of one window with another's — the reconcile cadence is + * what starves the download workers on large sets. + * @param idBufferBatches depth (in batches of [fetchBatch] ids) of the buffer + * between reconciliation and the download workers. Deeper keeps the workers fed + * across a window's round-trip gaps; memory is bounded by + * `idBufferBatches * fetchBatch` ids. + * + * **Subscription budget is the caller's job:** at peak this method holds + * `maxConcurrentReqs + reconcileConcurrency + 1` (keep-alive) subscriptions on + * the connection. Relays cap concurrent subscriptions per connection (NIP-11 + * `limitation.max_subscriptions`; e.g. strfry defaults to 20) and exceeding the + * cap can wedge the connection, not just fail the extra REQ — size the two knobs + * to fit the target relay. * @param onProgress optional `(needSoFar, downloaded)` ticks as work proceeds. * @param onEvent called once per distinct event, on the relay reader thread. */ +@OptIn(ExperimentalAtomicApi::class) suspend fun INostrClient.negentropySync( relay: NormalizedRelayUrl, filter: Filter, @@ -131,11 +152,13 @@ suspend fun INostrClient.negentropySync( maxConcurrentReqs: Int = 8, fetchBatch: Int = 500, idleTimeoutMs: Long = 120_000L, + reconcileConcurrency: Int = 1, + idBufferBatches: Int = maxConcurrentReqs * 4, onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, onEvent: (Event) -> Unit, ): NegentropySyncResult { - var need = 0 - var windows = 0 + val need = AtomicInt(0) + val windows = AtomicInt(0) var downloaded = 0 // Pin the relay in the pool's "desired" set for the whole sync. A NEG-OPEN is not @@ -156,17 +179,19 @@ suspend fun INostrClient.negentropySync( val producer = launch { try { - syncWindow( + syncPipeline( relay = relay, filter = filter, idleTimeoutMs = idleTimeoutMs, fetchBatch = fetchBatch, maxConcurrentReqs = maxConcurrentReqs, - onWindow = { windows++ }, + reconcileConcurrency = reconcileConcurrency, + idBufferBatches = idBufferBatches, + onWindow = { windows.incrementAndFetch() }, // Only accumulate here; progress is reported from the // single consumer loop below so the user callback is never // invoked from two coroutines at once. - onNeed = { need += it }, + onNeed = { need.addAndFetch(it) }, deliver = { events.send(it) }, ) } finally { @@ -177,7 +202,7 @@ suspend fun INostrClient.negentropySync( for (event in events) { downloaded++ onEvent(event) - onProgress?.invoke(need, downloaded) + onProgress?.invoke(need.load(), downloaded) if (maxEvents in 1..downloaded) break } @@ -190,10 +215,10 @@ suspend fun INostrClient.negentropySync( } return NegentropySyncResult( - needCount = need, + needCount = need.load(), haveCount = 0, downloaded = downloaded, - windows = windows, + windows = windows.load(), ) } @@ -204,6 +229,8 @@ suspend fun INostrClient.negentropySync( maxConcurrentReqs: Int = 8, fetchBatch: Int = 500, idleTimeoutMs: Long = 120_000L, + reconcileConcurrency: Int = 1, + idBufferBatches: Int = maxConcurrentReqs * 4, onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, onEvent: (Event) -> Unit, ): NegentropySyncResult = @@ -214,6 +241,8 @@ suspend fun INostrClient.negentropySync( maxConcurrentReqs = maxConcurrentReqs, fetchBatch = fetchBatch, idleTimeoutMs = idleTimeoutMs, + reconcileConcurrency = reconcileConcurrency, + idBufferBatches = idBufferBatches, onProgress = onProgress, onEvent = onEvent, ) @@ -257,6 +286,8 @@ suspend fun INostrClient.negentropySyncOrFetch( maxConcurrentReqs: Int = 8, fetchBatch: Int = 500, idleTimeoutMs: Long = 120_000L, + reconcileConcurrency: Int = 1, + idBufferBatches: Int = maxConcurrentReqs * 4, onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, onEvent: (Event) -> Unit, ): NegentropyOrFetchResult { @@ -283,6 +314,8 @@ suspend fun INostrClient.negentropySyncOrFetch( maxConcurrentReqs = maxConcurrentReqs, fetchBatch = fetchBatch, idleTimeoutMs = idleTimeoutMs, + reconcileConcurrency = reconcileConcurrency, + idBufferBatches = idBufferBatches, onProgress = onProgress, ) { accept(it) } NegentropyOrFetchResult(delivered, pagedFallback = false, negentropy = result, fallbackCause = null) @@ -306,6 +339,8 @@ suspend fun INostrClient.negentropySyncOrFetch( maxConcurrentReqs: Int = 8, fetchBatch: Int = 500, idleTimeoutMs: Long = 120_000L, + reconcileConcurrency: Int = 1, + idBufferBatches: Int = maxConcurrentReqs * 4, onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, onEvent: (Event) -> Unit, ): NegentropyOrFetchResult = @@ -316,60 +351,123 @@ suspend fun INostrClient.negentropySyncOrFetch( maxConcurrentReqs = maxConcurrentReqs, fetchBatch = fetchBatch, idleTimeoutMs = idleTimeoutMs, + reconcileConcurrency = reconcileConcurrency, + idBufferBatches = idBufferBatches, onProgress = onProgress, onEvent = onEvent, ) /** - * Recursively reconciles [filter] for [relay], splitting by `created_at` windows - * whenever the relay rejects the set as too large, and downloading the ids of each - * window as it resolves. Runs on a single coroutine; [deliver] funnels events out. + * The whole-sync pipeline: a single pool of [maxConcurrentReqs] download workers + * fed by up to [reconcileConcurrency] concurrent window reconciliations through a + * bounded buffer of [idBufferBatches] id batches. + * + * Why this shape (measured on a 2.6M-event production set — see + * `quartz/plans/2026-07-02-nostrclient-receiver-perf.md`): + * + * - The workers are GLOBAL, not per-window: with the old per-window pool the + * next window's reconcile round-trips only started after the previous + * window's last download joined, so the connection idled between windows. + * - Overflow-split windows go into a shared work queue processed by + * [reconcileConcurrency] reconcilers, overlapping their NEG round-trips. + * Each active reconcile holds one NEG session on the connection. + * - The id buffer decouples the bursty reconcile id stream from the steady + * download stream; it stays bounded so a slow consumer still back-pressures + * the relay (memory is O(idBufferBatches × fetchBatch) ids). * * Throws [NegentropySyncException] for any window negentropy cannot reconcile (a - * minimal window still over the cap, or an unavailable/erroring relay). + * minimal window still over the cap, or an unavailable/erroring relay); the + * failure cancels the whole pipeline. */ -private suspend fun INostrClient.syncWindow( +@OptIn(ExperimentalAtomicApi::class) +private suspend fun INostrClient.syncPipeline( relay: NormalizedRelayUrl, filter: Filter, idleTimeoutMs: Long, fetchBatch: Int, maxConcurrentReqs: Int, + reconcileConcurrency: Int, + idBufferBatches: Int, onWindow: () -> Unit, onNeed: (Int) -> Unit, deliver: suspend (Event) -> Unit, -) { - coroutineContext.ensureActive() +) = coroutineScope { + val idBatches = Channel>(idBufferBatches.coerceAtLeast(1)) - when (val outcome = downloadWindow(relay, filter, idleTimeoutMs, fetchBatch, maxConcurrentReqs, onNeed, deliver)) { - is ReconcileOutcome.Complete -> onWindow() - - is ReconcileOutcome.Overflow -> { - val lo = filter.since ?: 0L - val hi = filter.until ?: TimeUtils.now() - if (hi - lo <= MIN_WINDOW_SECONDS) { - // A minimal window that still overflows: negentropy genuinely can't - // enumerate this slice. Surface it — paging is the caller's call. - throw NegentropySyncException( - relay = relay, - window = filter, - reason = NegentropySyncException.Reason.OVER_MAX_SYNC_EVENTS, - detail = "created_at window [$lo, $hi] still exceeds the relay's max_sync_events", - ) - } else { - val mid = lo + (hi - lo) / 2 - syncWindow(relay, filter.copy(since = lo, until = mid), idleTimeoutMs, fetchBatch, maxConcurrentReqs, onWindow, onNeed, deliver) - syncWindow(relay, filter.copy(since = mid + 1, until = hi), idleTimeoutMs, fetchBatch, maxConcurrentReqs, onWindow, onNeed, deliver) + val workers = + List(maxConcurrentReqs.coerceAtLeast(1)) { + launch { + for (batch in idBatches) { + coroutineContext.ensureActive() + for (event in fetchByIds(relay, batch, idleTimeoutMs)) { + deliver(event) + } + } } } - is ReconcileOutcome.Failed -> - throw NegentropySyncException( - relay = relay, - window = filter, - reason = NegentropySyncException.Reason.UNAVAILABLE, - detail = outcome.detail, - ) - } + // Windows waiting for (or under) reconciliation. UNLIMITED so a reconciler + // re-queueing an overflow split never suspends while holding queue capacity + // (the split fan-out is tiny: two Filters per overflow). + val pending = Channel(Channel.UNLIMITED) + + // Queued-or-running windows. An overflow replaces one window with two + // (net +1); a completion is -1; the queue closes when it hits zero. + val remaining = AtomicInt(1) + pending.send(filter) + + val reconcilers = + List(reconcileConcurrency.coerceAtLeast(1)) { + launch { + for (window in pending) { + coroutineContext.ensureActive() + + val outcome = + reconcileStreaming(relay, window, idleTimeoutMs, fetchBatch, onNeed) { batch -> + idBatches.send(batch) + } + + when (outcome) { + is ReconcileOutcome.Complete -> { + onWindow() + if (remaining.decrementAndFetch() == 0) pending.close() + } + + is ReconcileOutcome.Overflow -> { + val lo = window.since ?: 0L + val hi = window.until ?: TimeUtils.now() + if (hi - lo <= MIN_WINDOW_SECONDS) { + // A minimal window that still overflows: negentropy + // genuinely can't enumerate this slice. Surface it — + // paging is the caller's call. + throw NegentropySyncException( + relay = relay, + window = window, + reason = NegentropySyncException.Reason.OVER_MAX_SYNC_EVENTS, + detail = "created_at window [$lo, $hi] still exceeds the relay's max_sync_events", + ) + } + val mid = lo + (hi - lo) / 2 + remaining.incrementAndFetch() + pending.send(window.copy(since = lo, until = mid)) + pending.send(window.copy(since = mid + 1, until = hi)) + } + + is ReconcileOutcome.Failed -> + throw NegentropySyncException( + relay = relay, + window = window, + reason = NegentropySyncException.Reason.UNAVAILABLE, + detail = outcome.detail, + ) + } + } + } + } + + reconcilers.joinAll() + idBatches.close() + workers.joinAll() } private sealed interface ReconcileOutcome { @@ -533,52 +631,6 @@ private fun isOverflow(reason: String): Boolean = reason.contains("too many", ignoreCase = true) || reason.startsWith("blocked", ignoreCase = true) -/** - * Reconciles [filter] and streams its ids straight into a bounded download pool, so - * reconciliation and download overlap and peak memory stays independent of the - * window's size. At most [maxConcurrentReqs] `REQ`s of [fetchBatch] ids are open at - * once; the id queue is bounded so a slow download back-pressures reconciliation. - * Returns the terminal [ReconcileOutcome]; events go out through [deliver]. - */ -private suspend fun INostrClient.downloadWindow( - relay: NormalizedRelayUrl, - filter: Filter, - idleTimeoutMs: Long, - fetchBatch: Int, - maxConcurrentReqs: Int, - onNeed: (Int) -> Unit, - deliver: suspend (Event) -> Unit, -): ReconcileOutcome = - coroutineScope { - val workerCount = maxConcurrentReqs.coerceAtLeast(1) - - // Bounded: when full, reconcileStreaming suspends instead of letting the - // relay's id stream accumulate. This is what keeps memory O(pipeline), not - // O(window). - val idBatches = Channel>(workerCount) - - val workers = - List(workerCount) { - launch { - for (batch in idBatches) { - coroutineContext.ensureActive() - for (event in fetchByIds(relay, batch, idleTimeoutMs)) { - deliver(event) - } - } - } - } - - val outcome = - reconcileStreaming(relay, filter, idleTimeoutMs, fetchBatch, onNeed) { batch -> - idBatches.send(batch) - } - - idBatches.close() - workers.joinAll() - outcome - } - /** * One `REQ` for [batch] ids; collects the matching events and returns them on * `EOSE`/close/timeout. All events for a single relay arrive on its one reader 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 f637c1d660..d4f6725dd2 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 @@ -39,6 +39,7 @@ import kotlinx.coroutines.async import kotlinx.coroutines.awaitAll import kotlinx.coroutines.coroutineScope import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.withTimeoutOrNull import okhttp3.OkHttpClient import okhttp3.Request import okhttp3.Response @@ -391,6 +392,96 @@ class BulkDownloadBenchmark { } } + // ------------------------------------------------------------------ + // negentropySync pipelining shootout: old sequential behavior + // (reconcileConcurrency=1, shallow buffer) vs the tuned pipeline, on the + // same production corpus, both capped at the same number of events. The + // per-connection by-id ceiling measured by ByIdFetchBenchmark (~5.5k/s on + // the 2.6M corpus) is the target; the old defaults measured ~1.6k/s. + // ------------------------------------------------------------------ + + private suspend fun negentropyVariant( + label: String, + httpClient: OkHttpClient, + maxEvents: Int, + maxConcurrentReqs: Int, + reconcileConcurrency: Int, + idBufferBatches: Int, + budgetMs: Long, + ) { + val relay = PROD_RELAY.normalizeRelayUrl() + val client = NostrClient(BasicOkHttpWebSocket.Builder { httpClient }) + val count = AtomicLong(0) + val start = System.nanoTime() + try { + val result = + withTimeoutOrNull(budgetMs) { + client.negentropySync( + relay = relay, + filter = Filter(kinds = listOf(PROD_KIND)), + maxEvents = maxEvents, + maxConcurrentReqs = maxConcurrentReqs, + reconcileConcurrency = reconcileConcurrency, + idBufferBatches = idBufferBatches, + ) { count.incrementAndGet() } + } + val wall = System.nanoTime() - start + report(label, count.get(), 0, wall) + if (result == null) { + println(" (budget ${budgetMs}ms exceeded — partial)") + } else { + println(" windows=${result.windows} need=${result.needCount}") + } + } catch (e: NegentropySyncException) { + val wall = System.nanoTime() - start + report(label, count.get(), 0, wall) + println(" !! failed: ${e.reason} ${e.message}") + } finally { + client.close() + } + } + + @Test + fun negentropyPipelineShootout() { + if (System.getenv("PROD_RELAY_BENCH") == null && System.getProperty("prodRelayBench") == null) { + println("negentropyPipelineShootout 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() + + println("=== NEGENTROPY PIPELINE SHOOTOUT: $PROD_RELAY kinds=[$PROD_KIND], cap 100k events ===") + runBlocking { + // Old behavior: sequential windows, shallow id buffer. + negentropyVariant( + "baseline (8 reqs, seq windows)", + httpClient, + maxEvents = 100_000, + maxConcurrentReqs = 8, + reconcileConcurrency = 1, + idBufferBatches = 8, + budgetMs = 240_000L, + ) + // Tuned: 12 + 4 + 1 keep-alive = 17 subs, under strfry's 20 cap. + negentropyVariant( + "tuned (12 reqs, 4 reconcilers)", + httpClient, + maxEvents = 100_000, + maxConcurrentReqs = 12, + reconcileConcurrency = 4, + idBufferBatches = 96, + budgetMs = 240_000L, + ) + } + 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