From abfea4e7b6841389c56bc966c8609d6d6ea5cae5 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 30 Jun 2026 21:48:11 +0000 Subject: [PATCH] feat(quartz): add high-level negentropy sync-and-download accessory MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Add `INostrClient.negentropySync` / `negentropySyncAsFlow`, a high-level NIP-77 accessory that downloads every event a relay holds matching a `Filter` and delivers each (deduped) through `onEvent` — mirroring the existing `fetchAllPages` accessory so downstream apps stop hand-rolling the `NegentropyManager` dance. It encapsulates the parts that make raw negentropy painful: - reconciles the relay's matched set (empty local set) via NegentropyManager - downloads the resulting ids through a bounded pool of concurrent REQs (`maxConcurrentReqs` subs of `fetchBatch` ids each, refilled on EOSE) - handles the relay-side cap (strfry `max_sync_events`, `NEG-ERR "blocked: too many query results"`) by splitting the filter into adaptive created_at windows; a minimal window that still can't reconcile (or a relay that doesn't speak NIP-77) falls back to `fetchAllPages` and reports it via `NegentropySyncResult.fellBackToPaging` - caps delivery at `maxEvents`, dedupes through a single consumer, and tears down all subscriptions + the neg session on completion/cancel Scope is controlled entirely by the `Filter` (per maintainer guidance the caller-supplied local-id delta interface is dropped in favour of a custom Filter), so the common call is one line. To drive NEG-OPEN on a single connection, add `INostrClient.getOrCreateRelay(url)` (default throws; NostrClient delegates to the pool). Because NEG-OPEN is a one-shot command that — unlike a REQ — is never replayed on reconnect, the accessory connects and waits for the relay to be ready before opening the session. Tests (quartz jvmAndroidTest, in-process relay): full download, maxEvents cap, clean teardown / no leaked subs, the Flow variant, and window-split + paging fallback against a relay that rejects the full reconcile. Co-Authored-By: Claude Opus 4.8 Claude-Session: https://claude.ai/code/session_01JmSyzdmKyiz3pPxUZ8Mg8Z --- .../nip01Core/relay/client/INostrClient.kt | 12 + .../nip01Core/relay/client/NostrClient.kt | 2 + .../NostrClientNegentropySyncAsFlowExt.kt | 84 ++++ .../NostrClientNegentropySyncExt.kt | 466 ++++++++++++++++++ .../relay/NostrClientNegentropySyncTest.kt | 166 +++++++ 5 files changed, 730 insertions(+) create mode 100644 quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncAsFlowExt.kt create mode 100644 quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt create mode 100644 quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/NostrClientNegentropySyncTest.kt diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt index 02b4ee74bc..125643864e 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt @@ -83,6 +83,18 @@ interface INostrClient : AutoCloseable { fun removeConnectionListener(listener: RelayConnectionListener) + /** + * Returns the [IRelayClient] for [url], creating and registering it in the + * connection pool if it is not there yet. + * + * Most callers should never need this — [subscribe]/[count]/[publish] manage + * the pool for you. It exists for accessories that must drive a single relay + * directly, such as NIP-77 negentropy (which sends `NEG-OPEN` and walks the + * reconciliation rounds on one connection). The default implementation throws; + * only a real pool-backed client can hand out relay clients. + */ + fun getOrCreateRelay(url: NormalizedRelayUrl): IRelayClient = throw UnsupportedOperationException("This INostrClient does not expose relay clients") + fun getReqFiltersOrNull(subId: String): Map>? fun getCountFiltersOrNull(subId: String): Map>? diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt index 693aa74661..f87b6622d4 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt @@ -342,6 +342,8 @@ class NostrClient( listeners.forEach { it.onCannotConnect(relay, errorMessage) } } + override fun getOrCreateRelay(url: NormalizedRelayUrl): IRelayClient = relayPool.getOrCreateRelay(url) + override fun addConnectionListener(listener: RelayConnectionListener) { listeners = listeners.plus(listener) } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncAsFlowExt.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncAsFlowExt.kt new file mode 100644 index 0000000000..298435ec1c --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncAsFlowExt.kt @@ -0,0 +1,84 @@ +/* + * 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.relay.client.INostrClient +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 kotlinx.coroutines.channels.awaitClose +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.callbackFlow + +/** + * Flow form of [negentropySync]: runs the sync while collected and emits the + * accumulated events as a growing list on each new arrival, then completes when + * the sync finishes (mirroring [com.vitorpamplona.quartz.nip01Core.relay.client.reqs.fetchAsFlow]). + * Cancelling the collector cancels the sync and tears down its subscriptions via + * [awaitClose]. + * + * See [negentropySync] for the meaning of every parameter. + */ +fun INostrClient.negentropySyncAsFlow( + relay: NormalizedRelayUrl, + filter: Filter, + maxEvents: Int = 0, + maxConcurrentReqs: Int = 8, + fetchBatch: Int = 500, + timeoutMs: Long = 30_000L, +): Flow> = + callbackFlow { + var current = listOf() + + negentropySync( + relay = relay, + filter = filter, + maxEvents = maxEvents, + maxConcurrentReqs = maxConcurrentReqs, + fetchBatch = fetchBatch, + timeoutMs = timeoutMs, + ) { event -> + current = current + event + trySend(current) + } + + close() + + awaitClose { } + } + +fun INostrClient.negentropySyncAsFlow( + relay: String, + filter: Filter, + maxEvents: Int = 0, + maxConcurrentReqs: Int = 8, + fetchBatch: Int = 500, + timeoutMs: Long = 30_000L, +): Flow> = + negentropySyncAsFlow( + relay = RelayUrlNormalizer.normalize(relay), + filter = filter, + maxEvents = maxEvents, + maxConcurrentReqs = maxConcurrentReqs, + fetchBatch = fetchBatch, + timeoutMs = timeoutMs, + ) 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 new file mode 100644 index 0000000000..6a9a092259 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt @@ -0,0 +1,466 @@ +/* + * 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.reqs.SubscriptionListener +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.normalizer.RelayUrlNormalizer +import com.vitorpamplona.quartz.nip77Negentropy.INegentropyListener +import com.vitorpamplona.quartz.nip77Negentropy.NegentropyManager +import com.vitorpamplona.quartz.utils.TimeUtils +import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.coroutineScope +import kotlinx.coroutines.ensureActive +import kotlinx.coroutines.flow.first +import kotlinx.coroutines.launch +import kotlinx.coroutines.withTimeoutOrNull +import kotlin.coroutines.coroutineContext +import kotlin.math.min + +/** + * Outcome of a [negentropySync] run. + * + * @property needCount ids the relay had that we lacked (i.e. everything that + * matched [Filter] on the relay — this sync always reconciles against an empty + * local set, so it downloads the full matched set). + * @property haveCount ids we had that the relay lacked. Always `0` here because + * the local set is empty; kept so the result mirrors a full NIP-77 reconcile. + * @property downloaded distinct events actually delivered through `onEvent`. + * @property windows number of `created_at` windows the matched set was split + * into (`1` when the relay reconciled the whole filter in one shot). + * @property fellBackToPaging `true` if at least one window could not be + * reconciled by negentropy (relay rejected it, e.g. strfry's + * `max_sync_events`, or did not support NIP-77) and was downloaded via + * [fetchAllPages] instead. + */ +class NegentropySyncResult( + val needCount: Int, + val haveCount: Int, + val downloaded: Int, + val windows: Int, + val fellBackToPaging: Boolean, +) + +/** + * Downloads every event a single [relay] holds matching [filter], delivering each + * one (deduped by id) through [onEvent]. A high-level wrapper over NIP-77 + * negentropy that hides the parts that make the raw protocol painful to use: + * + * 1. Reconciles the relay's matched set against an empty local set via a + * [NegentropyManager], accumulating the ids the relay has (`needIds`) across + * rounds until completion. + * 2. Downloads those ids through at most [maxConcurrentReqs] concurrent `REQ` + * subscriptions of [fetchBatch] ids each, refilling as each `EOSE` arrives, + * so a huge set never opens thousands of subs at once. + * 3. Handles the relay-side cap on negentropy (strfry's `max_sync_events`, + * observed as `NEG-ERR … "blocked: too many query results"`): the [filter] is + * split by `created_at` windows and each window reconciled on its own. A + * window that still overflows is halved and retried; a minimal window that + * still cannot reconcile (or a relay that does not speak NIP-77) falls back to + * [fetchAllPages] for that window and sets [NegentropySyncResult.fellBackToPaging]. + * + * Scope is controlled entirely by [filter] — narrow it (kinds, authors, `since`, + * tags, …) to download a slice instead of everything. The accessory keeps memory + * bounded by reconciling and downloading one window at a time and by capping the + * delivered set at [maxEvents]. + * + * Coroutine-cancellable: on completion, cancel, or reaching [maxEvents] all `REQ` + * subscriptions are unsubscribed and the negentropy session is closed and its + * listener removed, so nothing leaks. + * + * @param relay the relay to sync from. + * @param filter what to download. A single filter (NEG-OPEN is single-filter). + * @param maxEvents stop after delivering this many distinct events. `0` = unlimited. + * @param maxConcurrentReqs upper bound on simultaneously-open download `REQ`s. Keep + * it at or below the relay's per-connection subscription cap. + * @param fetchBatch ids per download `REQ`. + * @param timeoutMs max wait for a single reconcile round or download page's `EOSE`. + * @param onProgress optional `(needSoFar, downloaded)` ticks as work proceeds. + * @param onEvent called once per distinct event, on the relay reader thread. + */ +suspend fun INostrClient.negentropySync( + relay: NormalizedRelayUrl, + filter: Filter, + maxEvents: Int = 0, + maxConcurrentReqs: Int = 8, + fetchBatch: Int = 500, + timeoutMs: Long = 30_000L, + onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, + onEvent: (Event) -> Unit, +): NegentropySyncResult { + var need = 0 + var windows = 0 + var fellBack = false + + var downloaded = 0 + val seen = HashSet() + + coroutineScope { + // Single funnel for every delivered event (from REQ batches AND from the + // paging fallback), so dedup + the maxEvents cap + onEvent run on one + // coroutine even though the relay reader threads produce concurrently. + val events = Channel(Channel.UNLIMITED) + + val producer = + launch { + try { + syncWindow( + relay = relay, + filter = filter, + timeoutMs = timeoutMs, + fetchBatch = fetchBatch, + maxConcurrentReqs = maxConcurrentReqs, + onWindow = { windows++ }, + onPaged = { fellBack = true }, + // 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 }, + deliver = { events.trySend(it) }, + ) + } finally { + events.close() + } + } + + for (event in events) { + if (seen.add(event.id)) { + downloaded++ + onEvent(event) + onProgress?.invoke(need, downloaded) + if (maxEvents in 1..downloaded) break + } + } + + // If we broke out early (cap reached) the producer may still be working — + // stop it. If the producer finished normally this is a no-op. + producer.cancel() + } + + return NegentropySyncResult( + needCount = need, + haveCount = 0, + downloaded = downloaded, + windows = windows, + fellBackToPaging = fellBack, + ) +} + +suspend fun INostrClient.negentropySync( + relay: String, + filter: Filter, + maxEvents: Int = 0, + maxConcurrentReqs: Int = 8, + fetchBatch: Int = 500, + timeoutMs: Long = 30_000L, + onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, + onEvent: (Event) -> Unit, +): NegentropySyncResult = + negentropySync( + relay = RelayUrlNormalizer.normalize(relay), + filter = filter, + maxEvents = maxEvents, + maxConcurrentReqs = maxConcurrentReqs, + fetchBatch = fetchBatch, + timeoutMs = timeoutMs, + 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. + */ +private suspend fun INostrClient.syncWindow( + relay: NormalizedRelayUrl, + filter: Filter, + timeoutMs: Long, + fetchBatch: Int, + maxConcurrentReqs: Int, + onWindow: () -> Unit, + onPaged: () -> Unit, + onNeed: (Int) -> Unit, + deliver: (Event) -> Unit, +) { + coroutineContext.ensureActive() + + when (val outcome = reconcileWindow(relay, filter, timeoutMs)) { + is ReconcileOutcome.Ids -> { + onWindow() + onNeed(outcome.needIds.size) + downloadIds(relay, outcome.needIds, fetchBatch, maxConcurrentReqs, timeoutMs, deliver) + } + + 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 can't help — + // page it. Paging is bounded by created_at cursors, not the + // negentropy cap, so it always terminates. + onWindow() + onPaged() + fetchAllPages(relay, listOf(filter), timeoutMs, onEvent = deliver) + } else { + val mid = lo + (hi - lo) / 2 + syncWindow(relay, filter.copy(since = lo, until = mid), timeoutMs, fetchBatch, maxConcurrentReqs, onWindow, onPaged, onNeed, deliver) + syncWindow(relay, filter.copy(since = mid + 1, until = hi), timeoutMs, fetchBatch, maxConcurrentReqs, onWindow, onPaged, onNeed, deliver) + } + } + + is ReconcileOutcome.Failed -> { + // Relay doesn't speak NIP-77, disconnected, or timed out mid-reconcile. + // Fall back to paging for this window rather than giving up. + onWindow() + onPaged() + fetchAllPages(relay, listOf(filter), timeoutMs, onEvent = deliver) + } + } +} + +private sealed interface ReconcileOutcome { + /** Reconciliation completed; [needIds] are the ids the relay has that we lack. */ + class Ids( + val needIds: List, + ) : ReconcileOutcome + + /** Relay rejected the set as too large (strfry `max_sync_events`). */ + object Overflow : ReconcileOutcome + + /** Reconciliation could not complete (no NIP-77 support, disconnect, timeout). */ + object Failed : ReconcileOutcome +} + +/** + * Drives one NIP-77 reconciliation of [filter] against an EMPTY local set, so the + * resulting `needIds` are every id the relay holds for that filter. Registers a + * [NegentropyManager], sends `NEG-OPEN`, and walks the rounds until completion, + * error, or [timeoutMs]. Always closes the session and removes the listener. + */ +private suspend fun INostrClient.reconcileWindow( + relay: NormalizedRelayUrl, + filter: Filter, + timeoutMs: Long, +): ReconcileOutcome { + val relayClient = getOrCreateRelay(relay) + val subId = newSubId() + val signals = Channel(Channel.UNLIMITED) + val needIds = ArrayList() + + val listener = + object : INegentropyListener { + override fun onHaveIds( + relay: NormalizedRelayUrl, + subId: String, + haveIds: List, + ) { + // Empty local set: there is nothing the relay can lack. Ignore. + } + + override fun onNeedIds( + relay: NormalizedRelayUrl, + subId: String, + needIds: List, + ) { + signals.trySend(NegSignal.Need(needIds)) + } + + override fun onComplete( + relay: NormalizedRelayUrl, + subId: String, + ) { + signals.trySend(NegSignal.Complete) + } + + override fun onError( + relay: NormalizedRelayUrl, + subId: String, + reason: String, + ) { + signals.trySend(NegSignal.Error(reason)) + } + } + + val manager = NegentropyManager(listener) + addConnectionListener(manager) + try { + // NEG-OPEN is a one-shot command. Unlike a REQ — which the client replays + // from its active-request state every time a relay (re)connects — a dropped + // NEG-OPEN is never resent. `sendOrConnectAndSync` on a cold relay only + // kicks off the connect and silently drops the command, so we must connect + // and wait until the relay is ready before opening the session. + relayClient.connect() + val connected = + withTimeoutOrNull(timeoutMs) { + connectedRelaysFlow().first { relay in it } + } + if (connected == null) return ReconcileOutcome.Failed + + manager.startSync(relayClient, subId, filter, localEvents = emptyList()) + + val outcome = + withTimeoutOrNull(timeoutMs) { + while (true) { + when (val signal = signals.receive()) { + is NegSignal.Need -> needIds.addAll(signal.ids) + is NegSignal.Complete -> return@withTimeoutOrNull ReconcileOutcome.Ids(needIds) + is NegSignal.Error -> + return@withTimeoutOrNull if (isOverflow(signal.reason)) { + ReconcileOutcome.Overflow + } else { + ReconcileOutcome.Failed + } + } + } + @Suppress("UNREACHABLE_CODE") + ReconcileOutcome.Failed + } + + return outcome ?: ReconcileOutcome.Failed + } finally { + manager.closeSync(relayClient, subId) + removeConnectionListener(manager) + signals.close() + } +} + +private sealed interface NegSignal { + class Need( + val ids: List, + ) : NegSignal + + object Complete : NegSignal + + class Error( + val reason: String, + ) : NegSignal +} + +/** + * strfry sends `["NEG-ERR", subId, "blocked: too many query results"]` when a + * NEG-OPEN matches more than `relay__negentropy__maxSyncEvents`. Match that + * verbatim, plus a looser contains-check so equivalent wording from other relays + * still triggers the window split rather than aborting. + */ +private fun isOverflow(reason: String): Boolean = + reason == "blocked: too many query results" || + reason.contains("too many", ignoreCase = true) || + reason.startsWith("blocked", ignoreCase = true) + +/** + * Downloads [ids] from [relay] through at most [maxConcurrentReqs] concurrent + * `REQ`s of [fetchBatch] ids each. A fixed worker pool drains a batch queue, so at + * most [maxConcurrentReqs] subscriptions are ever open at once, each refilled as + * its `EOSE` arrives. + */ +private suspend fun INostrClient.downloadIds( + relay: NormalizedRelayUrl, + ids: List, + fetchBatch: Int, + maxConcurrentReqs: Int, + timeoutMs: Long, + deliver: (Event) -> Unit, +) { + if (ids.isEmpty()) return + + val batches = Channel>(Channel.UNLIMITED) + val chunks = ids.chunked(fetchBatch) + chunks.forEach { batches.trySend(it) } + batches.close() + + val workers = min(maxConcurrentReqs.coerceAtLeast(1), chunks.size) + + coroutineScope { + repeat(workers) { + launch { + for (batch in batches) { + coroutineContext.ensureActive() + fetchByIds(relay, batch, timeoutMs, deliver) + } + } + } + } +} + +/** One `REQ` for [batch] ids; delivers each event, returns on `EOSE`/close/timeout. */ +private suspend fun INostrClient.fetchByIds( + relay: NormalizedRelayUrl, + batch: List, + timeoutMs: Long, + deliver: (Event) -> Unit, +) { + val subId = newSubId() + val done = Channel(Channel.CONFLATED) + + val listener = + object : SubscriptionListener { + override fun onEvent( + event: Event, + isLive: Boolean, + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + deliver(event) + } + + override fun onEose( + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + done.trySend(Unit) + } + + override fun onClosed( + message: String, + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + done.trySend(Unit) + } + + override fun onCannotConnect( + relay: NormalizedRelayUrl, + message: String, + forFilters: List?, + ) { + done.trySend(Unit) + } + } + + try { + subscribe(subId, mapOf(relay to listOf(Filter(ids = batch))), listener) + withTimeoutOrNull(timeoutMs) { + done.receive() + } + } finally { + unsubscribe(subId) + done.close() + } +} + +/** Seconds: a window this small that still overflows is paged instead of split. */ +private const val MIN_WINDOW_SECONDS = 1L diff --git a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/NostrClientNegentropySyncTest.kt b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/NostrClientNegentropySyncTest.kt new file mode 100644 index 0000000000..20e573244c --- /dev/null +++ b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/NostrClientNegentropySyncTest.kt @@ -0,0 +1,166 @@ +/* + * 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.InProcessRelays +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.negentropySync +import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.negentropySyncAsFlow +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer +import com.vitorpamplona.quartz.nip77Negentropy.NegentropySettings +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.cancel +import kotlinx.coroutines.flow.last +import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.withTimeout +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFalse +import kotlin.test.assertTrue + +class NostrClientNegentropySyncTest : RelayClientTest() { + @Test + fun fullDownloadDeliversEveryEvent() = + runBlocking { + defaultRelay.preload(SyntheticEvents.batch(20, kind = 1)) + + val got = mutableListOf() + val result = + withTimeout(20_000) { + client.negentropySync( + relay = defaultRelayUrl, + filter = Filter(kinds = listOf(1)), + ) { got.add(it) } + } + + assertEquals(20, got.size, "every event should be delivered") + assertEquals(20, got.map { it.id }.toSet().size, "no duplicates") + assertEquals(20, result.needCount) + assertEquals(0, result.haveCount) + assertEquals(20, result.downloaded) + assertEquals(1, result.windows, "small set reconciles in a single window") + assertFalse(result.fellBackToPaging) + } + + @Test + fun maxEventsCapsDelivery() = + runBlocking { + defaultRelay.preload(SyntheticEvents.batch(20, kind = 1)) + + val got = mutableListOf() + val result = + withTimeout(20_000) { + client.negentropySync( + relay = defaultRelayUrl, + filter = Filter(kinds = listOf(1)), + maxEvents = 10, + fetchBatch = 5, + ) { got.add(it) } + } + + assertEquals(10, got.size, "delivery stops at maxEvents") + assertEquals(10, result.downloaded) + } + + @Test + fun cleanTeardownLeavesNoSubscriptions() = + runBlocking { + defaultRelay.preload(SyntheticEvents.batch(15, kind = 1)) + + withTimeout(20_000) { + client.negentropySync( + relay = defaultRelayUrl, + filter = Filter(kinds = listOf(1)), + fetchBatch = 4, + ) { } + } + + assertTrue( + client.activeRequests(defaultRelayUrl).isEmpty(), + "all download subscriptions must be closed after the sync", + ) + } + + @Test + fun flowVariantEmitsAllEvents() = + runBlocking { + defaultRelay.preload(SyntheticEvents.batch(12, kind = 1)) + + val last = + withTimeout(20_000) { + client + .negentropySyncAsFlow( + relay = defaultRelayUrl, + filter = Filter(kinds = listOf(1)), + ).last() + } + + assertEquals(12, last.size) + assertEquals(12, last.map { it.id }.toSet().size) + } + + /** + * A relay that caps negentropy below the matched-set size (strfry's + * `max_sync_events`) and whose events all share one `created_at`, so no + * `created_at` window can separate them. The sync must still download every + * event — by splitting down to a minimal window and then falling back to + * [com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAllPages]. + */ + @Test + fun windowSplitAndPagingFallbackOnOverflow() = + runBlocking { + val hub = InProcessRelays(negentropySettings = NegentropySettings(maxSyncEvents = 3)) + val scope = CoroutineScope(Dispatchers.Default + SupervisorJob()) + val client = NostrClient(hub, scope) + try { + val url = RelayUrlNormalizer.normalize("ws://127.0.0.1:7780/") + // 10 events, all at the same created_at: created_at windowing can + // never split them, so the minimal window still overflows the cap. + val events = (1..10).map { SyntheticEvents.fakeEvent(idSeed = it, kind = 1, createdAt = 1000L) } + hub.getOrCreate(url).preload(events) + + val got = mutableListOf() + val result = + withTimeout(60_000) { + client.negentropySync( + relay = url, + filter = Filter(kinds = listOf(1)), + ) { got.add(it) } + } + + assertEquals(10, got.map { it.id }.toSet().size, "all events downloaded despite the cap") + assertEquals(10, result.downloaded) + assertTrue(result.fellBackToPaging, "the over-cap minimal window should page") + assertTrue(result.windows > 1, "the set must be split into multiple created_at windows") + } finally { + client.disconnect() + scope.cancel() + hub.close() + } + } +}