feat(cli): add Context.drainAllPages + shared fetchAllPagesFromPool accessory

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 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_015YEbdqCRPkszkGCoi89RMt
This commit is contained in:
Claude
2026-07-06 20:20:57 +00:00
parent 54c8cefe69
commit 590731f356
6 changed files with 249 additions and 64 deletions
@@ -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<NormalizedRelayUrl>,
filters: Map<NormalizedRelayUrl, List<Filter>>,
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<Filter>,
onNewPage: (Long) -> Unit,
onEvent: (Event) -> Unit,
): Int = fetchAllPages(relay, filters, RELAY_TIMEOUT_MS, onNewPage, onEvent)
}
@@ -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<NormalizedRelayUrl>,
@@ -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<NormalizedRelayUrl, List<Filter>>,
timeoutMs: Long = 30_000,
maxConcurrentRelays: Int = 8,
): List<Pair<NormalizedRelayUrl, Event>> {
if (filters.isEmpty()) return emptyList()
val collected = mutableListOf<Pair<NormalizedRelayUrl, Event>>()
// 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<Pair<NormalizedRelayUrl, Event>>(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
@@ -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 <nevent1…|naddr1…|nprofile1…|npub1…|note1…|name@domain>
*
@@ -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()
@@ -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<NormalizedRelayUrl, List<Filter>>,
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()
}
}
}
}
}
@@ -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<Pair<NormalizedRelayUrl, Event>>())
val completed = ConcurrentHashMap<NormalizedRelayUrl, Int>()
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])
}
}