From 590731f3561a4a079b7ec07aa7cbb529598d3b0c Mon Sep 17 00:00:00 2001 From: Claude Date: Mon, 6 Jul 2026 20:20:57 +0000 Subject: [PATCH] feat(cli): add Context.drainAllPages + shared fetchAllPagesFromPool accessory MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Amy's one-shot queries all go through Context.drain, a single REQ drained to EOSE — so a relay that caps its REQ response (strfry's per-REQ limit, ~500) silently truncates the result with no way to page past it. Extract the per-relay fetchAllPages fan-out that already lived privately in EventSync into a reusable quartz accessory, fetchAllPagesFromPool: a sliding-window pool (maxConcurrentRelays) that paginates each relay on its own `until` cursor, tags every event with its source relay, and does not dedup across relays. EventSync now delegates to it (its private downloadPool/ downloadFromRelay are deleted — no behavior change: perRelayFilters is already ordered by and complete over the relay list). Add Context.drainAllPages, the paged sibling of drain: same verify+store and per-relay tagging, but fully draining sets larger than one REQ. Wire it into `amy fetch` behind --paginate/--all (filter mode only), pushing the limit into the filter so paging stays bounded. sync (NIP-77) and fetch stay separate interfaces. Tests: fetchAllPagesFromPool fan-out/tagging/no-cross-relay-dedup. Co-Authored-By: Claude Opus 4.8 Claude-Session: https://claude.ai/code/session_015YEbdqCRPkszkGCoi89RMt --- .../loggedIn/relays/eventsync/EventSync.kt | 65 +------------ .../wallet/wizard/CashuWalletDiscovery.kt | 5 +- .../com/vitorpamplona/amethyst/cli/Context.kt | 50 ++++++++++ .../amethyst/cli/commands/FetchCommand.kt | 20 +++- .../NostrClientFetchAllPagesPoolExt.kt | 93 +++++++++++++++++++ .../relay/NostrClientFetchAllPagesPoolTest.kt | 80 ++++++++++++++++ 6 files changed, 249 insertions(+), 64 deletions(-) create mode 100644 quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientFetchAllPagesPoolExt.kt create mode 100644 quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/NostrClientFetchAllPagesPoolTest.kt diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/relays/eventsync/EventSync.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/relays/eventsync/EventSync.kt index e7c962442d..d589c6de78 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/relays/eventsync/EventSync.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/relays/eventsync/EventSync.kt @@ -27,7 +27,7 @@ import com.vitorpamplona.amethyst.ui.screen.loggedIn.relays.eventsync.EventSync. 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.accessories.fetchAllPages +import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAllPagesFromPool import com.vitorpamplona.quartz.nip01Core.relay.client.listeners.RelayConnectionListener import com.vitorpamplona.quartz.nip01Core.relay.client.single.IRelayClient import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message @@ -46,10 +46,7 @@ import kotlinx.coroutines.Job import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.update -import kotlinx.coroutines.isActive import kotlinx.coroutines.launch -import kotlinx.coroutines.supervisorScope -import kotlinx.coroutines.sync.Semaphore import java.util.concurrent.ConcurrentHashMap import kotlin.coroutines.cancellation.CancellationException @@ -503,9 +500,10 @@ class EventSync( clientBuilder().use { client -> client.addConnectionListener(okListener) try { - client.downloadFromPool( - relays = relaysToProcess, + client.fetchAllPagesFromPool( filters = perRelayFilters, + timeoutMs = RELAY_TIMEOUT_MS, + maxConcurrentRelays = MAX_CONCURRENT_RELAYS, onNewPage = { until, sourceRelay -> _liveActivity.value.runningRelays[sourceRelay] ?.pageUntil @@ -564,7 +562,7 @@ class EventSync( ) } }, - onRelayComplete = { relay -> + onRelayComplete = { relay, _ -> _liveActivity.update { val newCompleted = it.runningRelays[relay] it.copy( @@ -615,57 +613,4 @@ class EventSync( } } } - - /** - * Maintains a sliding window of up to [MAX_CONCURRENT_RELAYS] active relay workers. - * As soon as one relay finishes (all pages exhausted), the next relay from [relays] - * starts immediately — no waiting for an entire batch to drain. - * - * [onEvent] receives the event and the URL of the relay it came from. - */ - private suspend fun INostrClient.downloadFromPool( - relays: List, - filters: Map>, - onNewPage: (Long, NormalizedRelayUrl) -> Unit, - onEvent: (Event, NormalizedRelayUrl) -> Unit, - onRelayStart: (NormalizedRelayUrl) -> Unit, - onRelayComplete: (NormalizedRelayUrl) -> Unit, - ) { - val semaphore = Semaphore(MAX_CONCURRENT_RELAYS) - supervisorScope { - for (relay in relays) { - if (!isActive) break - semaphore.acquire() - launch { - try { - onRelayStart(relay) - filters[relay]?.let { filtersForRelay -> - downloadFromRelay( - relay = relay, - filters = filtersForRelay, - onNewPage = { onNewPage(it, relay) }, - onEvent = { onEvent(it, relay) }, - ) - } ?: 0 - onRelayComplete(relay) - } finally { - semaphore.release() - } - } - } - } - } - - /** - * Fetches all pages from a single [relay] using paginated `until` cursors. - * Delegates to the Quartz [downloadFromRelay] extension. - * - * @return total number of events received across all pages. - */ - private suspend fun INostrClient.downloadFromRelay( - relay: NormalizedRelayUrl, - filters: List, - onNewPage: (Long) -> Unit, - onEvent: (Event) -> Unit, - ): Int = fetchAllPages(relay, filters, RELAY_TIMEOUT_MS, onNewPage, onEvent) } diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/wallet/wizard/CashuWalletDiscovery.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/wallet/wizard/CashuWalletDiscovery.kt index 45ad8720f2..a68f4e356b 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/wallet/wizard/CashuWalletDiscovery.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/wallet/wizard/CashuWalletDiscovery.kt @@ -177,8 +177,9 @@ class CashuWalletDiscovery( /** * Sliding-window relay crawl: keeps up to [MAX_CONCURRENT_RELAYS] relays - * paginating at once, starting the next as soon as one finishes. Mirrors - * EventSync.downloadFromPool but only collects (no republish). + * paginating at once, starting the next as soon as one finishes. Same shape as + * the shared `fetchAllPagesFromPool` accessory, but with one filter list for + * every relay and collect-only (no per-relay tagging, no republish). */ private suspend fun INostrClient.crawlPool( relays: List, diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt index dffb600206..f62e624a1e 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt @@ -42,6 +42,7 @@ import com.vitorpamplona.quartz.nip01Core.crypto.verify import com.vitorpamplona.quartz.nip01Core.jackson.JacksonMapper import com.vitorpamplona.quartz.nip01Core.metadata.MetadataEvent import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient +import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAllPagesFromPool import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.publishAndConfirmDetailed import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener import com.vitorpamplona.quartz.nip01Core.relay.client.single.newSubId @@ -72,6 +73,8 @@ import com.vitorpamplona.quartz.nip87Ecash.recommendation.MintRecommendationEven import kotlinx.coroutines.CompletableDeferred import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.channels.Channel.Factory.UNLIMITED +import kotlinx.coroutines.coroutineScope +import kotlinx.coroutines.launch import kotlinx.coroutines.selects.select import kotlinx.coroutines.withTimeoutOrNull import okhttp3.OkHttpClient @@ -479,6 +482,53 @@ class Context( return collected } + /** + * Like [drain], but paginates every relay to completion via + * [fetchAllPagesFromPool] instead of stopping at the first EOSE — so a query + * larger than a relay's per-`REQ` cap (strfry's `limit`, ~500) is fully + * retrieved instead of silently truncated. Each relay is walked on its own + * `until` cursor, up to [maxConcurrentRelays] at once, and every event still + * funnels through [verifyAndStore]; the result is tagged by relay exactly like + * [drain] (and, like [drain], is NOT deduped across relays — callers dedup by + * id). + * + * Bound the work with the filters' `limit`: each relay pages until it reaches + * the limit, so an unbounded filter pages that relay's entire matching history. + * A `search` filter is fetched as a single relevance-ranked page (see + * [com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAllPages]). + */ + suspend fun drainAllPages( + filters: Map>, + timeoutMs: Long = 30_000, + maxConcurrentRelays: Int = 8, + ): List> { + if (filters.isEmpty()) return emptyList() + val collected = mutableListOf>() + // fetchAllPages' onEvent can't suspend, but verifyAndStore does — bridge + // through a channel and verify+store single-threaded in one consumer so the + // store writes stay serialized (same shape as `drain`). + val eventChannel = Channel>(UNLIMITED) + coroutineScope { + val consumer = + launch { + for ((relay, event) in eventChannel) { + if (verifyAndStore(event)) collected.add(relay to event) + } + } + try { + client.fetchAllPagesFromPool( + filters = filters, + timeoutMs = timeoutMs, + maxConcurrentRelays = maxConcurrentRelays, + ) { event, relay -> eventChannel.trySend(relay to event) } + } finally { + eventChannel.close() + } + consumer.join() + } + return collected + } + /** * Publish [request] to [relays], then wait for the FIRST event matching [responseFilter] * — a live reply that arrives after our own EOSE, which [drain] would miss (it returns at diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/FetchCommand.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/FetchCommand.kt index 2adffc664c..109b347f26 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/FetchCommand.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/FetchCommand.kt @@ -36,7 +36,8 @@ import com.vitorpamplona.quartz.nip65RelayList.AdvertisedRelayListEvent /** * `amy fetch [--kind …] [--author …] [--id …] [--tag …] [--since/--until TS] - * [--limit N] [--search TEXT] [--relay URL[,URL…]] [--timeout SECS]` + * [--limit N] [--search TEXT] [--relay URL[,URL…]] [--timeout SECS] + * [--paginate]` * * amy fetch * @@ -53,6 +54,12 @@ import com.vitorpamplona.quartz.nip65RelayList.AdvertisedRelayListEvent * * Results are deduplicated by id, sorted newest-first, capped at `--limit` * (default 100), and emitted as full event JSON under an `events` array. + * + * By default (filter mode) a single `REQ` is drained to EOSE, so a relay that + * caps its response (strfry's per-`REQ` `limit`, ~500) truncates the result. + * `--paginate` (alias `--all`) instead walks each relay page-by-page on `until` + * cursors up to `--limit`, fully draining sets larger than one `REQ` — the + * multi-relay [Context.drainAllPages] path. Code mode is always single-shot. */ object FetchCommand { suspend fun run( @@ -71,13 +78,22 @@ object FetchCommand { } val filter = RawEventSupport.buildFilter(args) + // --paginate/--all walks each relay past its per-REQ cap (strfry's ~500) + // by following `until` cursors, bounded by `limit`; default stops at the + // first EOSE like nak's `req`. + val paginate = args.bool("paginate") || args.bool("all") Context.open(dataDir).use { ctx -> ctx.prepare() val relays = RawEventSupport.queryTargets(ctx, args) if (relays.isEmpty()) return Output.error("no_relays", "no relays available; pass --relay or run `amy relay add`") - val received = ctx.drain(relays.associateWith { listOf(filter) }, timeoutMs) + val received = + if (paginate) { + ctx.drainAllPages(relays.associateWith { listOf(filter.copy(limit = limit)) }, timeoutMs) + } else { + ctx.drain(relays.associateWith { listOf(filter) }, timeoutMs) + } val events = received .asSequence() diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientFetchAllPagesPoolExt.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientFetchAllPagesPoolExt.kt new file mode 100644 index 0000000000..c6d4ffeca2 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientFetchAllPagesPoolExt.kt @@ -0,0 +1,93 @@ +/* + * 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 kotlinx.coroutines.isActive +import kotlinx.coroutines.launch +import kotlinx.coroutines.supervisorScope +import kotlinx.coroutines.sync.Semaphore + +/** + * Fans [fetchAllPages] out across every relay in [filters] — the multi-relay form + * of the single-relay downloader. Each relay is paginated independently on its own + * `until` cursor, with at most [maxConcurrentRelays] relays in flight at once: as + * soon as one drains (all pages exhausted) the next from [filters] starts, so the + * concurrency window stays full instead of waiting for a whole batch to finish. + * + * Every delivered event is tagged with the relay it came from. **No cross-relay + * dedup is done here** — the same event id can arrive from several relays, exactly + * like a fan-out `REQ`; dedup downstream if you need a distinct set. [onEvent] runs + * on the delivering relay's reader thread and must not suspend (bridge through a + * channel if your sink suspends). + * + * Relays are isolated by [supervisorScope]: one relay throwing — or all its pages + * timing out — fails only that relay's branch, never the others. A relay that can't + * connect simply yields zero events (its first page EOSEs empty) and completes + * normally. Cancelling the caller cancels every branch. + * + * @param filters per-relay filter lists; the key set is the relays queried, in + * iteration order (pass a [LinkedHashMap]/`associateWith` result to control it). + * A `search` filter is fetched as a single relevance page — see [fetchAllPages]. + * @param timeoutMs per-page EOSE timeout handed to each relay's [fetchAllPages]. + * @param maxConcurrentRelays upper bound on relays paginating at once (≥ 1). + * @param onNewPage optional `(until, relay)` tick before each non-first page. + * @param onRelayStart optional hook fired as each relay's download begins. + * @param onRelayComplete optional `(relay, totalEvents)` hook fired when a relay + * drains (or errors out to an empty first page). + * @param onEvent called once per delivered event with its source relay. + */ +suspend fun INostrClient.fetchAllPagesFromPool( + filters: Map>, + timeoutMs: Long = 30_000L, + maxConcurrentRelays: Int = 8, + onNewPage: ((until: Long, relay: NormalizedRelayUrl) -> Unit)? = null, + onRelayStart: ((relay: NormalizedRelayUrl) -> Unit)? = null, + onRelayComplete: ((relay: NormalizedRelayUrl, totalEvents: Int) -> Unit)? = null, + onEvent: (event: Event, relay: NormalizedRelayUrl) -> Unit, +) { + if (filters.isEmpty()) return + val semaphore = Semaphore(maxConcurrentRelays.coerceAtLeast(1)) + supervisorScope { + for ((relay, filtersForRelay) in filters) { + if (!isActive) break + semaphore.acquire() + launch { + try { + onRelayStart?.invoke(relay) + val total = + fetchAllPages( + relay = relay, + filters = filtersForRelay, + timeoutMs = timeoutMs, + onNewPage = onNewPage?.let { cb -> { until -> cb(until, relay) } }, + ) { event -> onEvent(event, relay) } + onRelayComplete?.invoke(relay, total) + } finally { + semaphore.release() + } + } + } + } +} diff --git a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/NostrClientFetchAllPagesPoolTest.kt b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/NostrClientFetchAllPagesPoolTest.kt new file mode 100644 index 0000000000..a3c962882e --- /dev/null +++ b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/NostrClientFetchAllPagesPoolTest.kt @@ -0,0 +1,80 @@ +/* + * 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.accessories.fetchAllPagesFromPool +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.runBlocking +import java.util.Collections +import java.util.concurrent.ConcurrentHashMap +import kotlin.test.Test +import kotlin.test.assertEquals + +class NostrClientFetchAllPagesPoolTest : RelayClientTest() { + /** + * The pool fans [fetchAllPagesFromPool] across relays, tags each event with the + * relay it came from, and — like a fan-out REQ — does NOT dedup across relays: + * an event held by two relays is delivered twice, once per source. + */ + @Test + fun fansOutPerRelayTagsSourceAndDoesNotDedupAcrossRelays() = + runBlocking { + val relayAUrl = RelayUrlNormalizer.normalize("ws://relay-a/") + val relayBUrl = RelayUrlNormalizer.normalize("ws://relay-b/") + + // One shared event lives on BOTH relays; the rest are relay-exclusive. + val shared = SyntheticEvents.fakeEvent(idSeed = 1) + val onlyA = (2..4).map { SyntheticEvents.fakeEvent(idSeed = it) } + val onlyB = (5..9).map { SyntheticEvents.fakeEvent(idSeed = it) } + hub.getOrCreate(relayAUrl).preload(onlyA + shared) + hub.getOrCreate(relayBUrl).preload(onlyB + shared) + + // onEvent runs on each relay's reader thread → collect thread-safely. + val received = Collections.synchronizedList(mutableListOf>()) + val completed = ConcurrentHashMap() + + client.fetchAllPagesFromPool( + filters = + linkedMapOf( + relayAUrl to listOf(Filter(kinds = listOf(1))), + relayBUrl to listOf(Filter(kinds = listOf(1))), + ), + onRelayComplete = { relay, total -> completed[relay] = total }, + ) { event, relay -> received.add(relay to event) } + + val fromA = received.filter { it.first == relayAUrl } + val fromB = received.filter { it.first == relayBUrl } + // A: 3 exclusive + shared = 4 ; B: 5 exclusive + shared = 6. + assertEquals(4, fromA.size, "every event must be tagged with relay A") + assertEquals(6, fromB.size, "every event must be tagged with relay B") + // The shared id arrives once per relay — not deduped across relays. + assertEquals(2, received.count { it.second.id == shared.id }, "shared event must arrive from both relays") + // onRelayComplete reports each relay's fetched total. + assertEquals(4, completed[relayAUrl]) + assertEquals(6, completed[relayBUrl]) + } +}