From 22b8089a96171c0fed06110f73d707329567093d Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 8 Jul 2026 17:35:52 +0000 Subject: [PATCH] Revert "fix(graperank): paginate aggregator recovery so the page cap can't truncate it" This reverts commit f19b8052b0eab841909052223fdab019dc0f11f8. --- .../graperank/GrapeRankDataCrawler.kt | 57 +++++++------------ 1 file changed, 21 insertions(+), 36 deletions(-) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankDataCrawler.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankDataCrawler.kt index 5cc896eef0..29dc9600d8 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankDataCrawler.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankDataCrawler.kt @@ -27,7 +27,6 @@ import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.AdaptiveRelayLimiter import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.DrainFailure import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.classifyDrainFailure -import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAllPages 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 @@ -686,18 +685,19 @@ class GrapeRankDataCrawler( val before = contactListsFed log("[graperank] aggregator recovery: ${stragglers.size} stragglers via ${aggregators.size} aggregators") - // Chunk the stragglers once. Each chunk is queried on its OWN request: - // the big indexers return nothing for a filter carrying the whole set, so - // AUTHORS_PER_FILTER-sized chunks keep every request answerable. Ask ONLY - // for kind:3 — the contact list we're missing. A multi-kind filter is - // useless against these indexers: user.kindpag.es caps its response at ~100 - // events per page (it ignores our limit), so a kinds=[3,10000,1984,10002] - // query comes back 100× kind:10002 and 0× kind:3 — the abundant relay lists - // crowd the contact lists out. Their kind:10002 is already fetched in bulk - // by [ensureRelayLists]; mutes/reports still come from the outbox model. - val chunks = - stragglers.chunked(AUTHORS_PER_FILTER).map { chunk -> - Filter(kinds = listOf(ContactListEvent.KIND), authors = chunk) + // Build the query against the full straggler set BEFORE any folding (the + // filter lists are materialized here, so later mutation of `stragglers` is + // safe). Ask ONLY for kind:3 — the contact list we're missing. A multi-kind + // filter is useless against the big indexers: user.kindpag.es caps its + // response at ~100 events per REQ (it ignores our limit), so a + // kinds=[3,10000,1984,10002] query comes back 100× kind:10002 and 0× + // kind:3 — the abundant relay lists crowd the contact lists out entirely. + // Asked for kind:3 alone it returns them in a few seconds. Their kind:10002 + // is already fetched in bulk by [ensureRelayLists]; mutes/reports still come + // from the outbox model. The aggregator's job here is only the lists. + val filters = + aggregators.associateWith { + stragglers.chunked(AUTHORS_PER_FILTER).map { chunk -> Filter(kinds = listOf(ContactListEvent.KIND), authors = chunk) } } // Fold one delivered contact list per straggler, exactly once. @@ -715,30 +715,15 @@ class GrapeRankDataCrawler( } relaysContacted += aggregators - // Paginate every aggregator with `until` cursors instead of the crawl's - // single-shot [drainGated]: that grabs one page and stops, so against a - // relay that caps a page at ~100 events any straggler beyond the newest 100 - // is silently dropped. [fetchAllPages] walks the whole filter to exhaustion. - // Relays run concurrently; each relay's chunks run sequentially so only one - // subscription is live per connection (staying under per-relay sub limits), - // gated by the [limiter] like every other query. Events land on a channel - // off the reader threads, then are verified/persisted and folded once. - val sink = Channel>(Channel.UNLIMITED) - coroutineScope { - for (agg in aggregators) { - launch { - for (chunk in chunks) { - limiter.withPermit(agg) { - client.fetchAllPages(agg, listOf(chunk), config.parkTimeoutMs) { ev -> - sink.trySend(agg to ev) - } - } - } - } - } + // Fast deliveries fold immediately; a slow aggregator parks and its late + // kind:3 arrives on [lateHarvest], which we drain until the parked units + // finish — so a big aggregator that can't answer within the fast window is + // still fully harvested here instead of being abandoned. + foldAgg(drainGated(filters, null)) + while (parkedInFlight.load() > 0L) { + withTimeoutOrNull(PARK_POLL_MS) { lateHarvest.receive() }?.let { foldAgg(listOf(it)) } } - sink.close() - foldAgg(persist(buildList { for (e in sink) add(e) })) + while (true) foldAgg(listOf(lateHarvest.tryReceive().getOrNull() ?: break)) log("[graperank] aggregator recovery: +${contactListsFed - before} contact lists") }