From 667358a06a40c0f36828b3b96b0652aa28598c2a Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 7 Jul 2026 15:41:31 +0000 Subject: [PATCH] refactor(quartz): extract GrapeRankDataCrawler to commonMain The web-of-trust crawl (~400 lines: outbox routing, sharded backbone sweep, Phase-B worker pool, relay-list discovery, report-deletion fetch, warm pool) was making the CLI's GrapeRankCommand unmaintainably large. Move it into a reusable, KMP-portable GrapeRankDataCrawler in quartz commonMain. The crawler takes a NostrClient + IEventStore + AdaptiveRelayLimiter, injected relay policy (discovery + content-fallback sets, since those defaults live in app code, not the protocol library), and a log callback; it streams contact lists into a TrustGraphBuilder and returns crawl Stats. GrapeRankCommand shrinks to arg-parsing + offline load + scoring + publish + sub-verbs, delegating the online path to the crawler. To reach commonMain (portable to every target, incl. iOS): - Add ConcurrentMap / ConcurrentSet expect classes under utils/concurrent, with jvmAndroid actuals (java.util.concurrent) and native actuals (copy-on-write over kotlin.concurrent.atomics.AtomicReference, mirroring ConcurrentHashCache). commonMain has no ConcurrentHashMap, and the crawl's producer/consumer/drain- worker state needs atomic getOrPut/merge plus a concurrent set. - Move AdaptiveRelayLimiter and DrainFailure/classifyDrainFailure from cli to quartz commonMain (java atomics -> kotlin.concurrent.atomics, ConcurrentHashMap -> ConcurrentMap, System.currentTimeMillis -> TimeUtils.nowMillis, stderr -> Log). - The gated drain (REQ-size splitting, per-relay permits, verify+store) moves into the crawler; Context.drain loses its now-unused gatePerRelay path. Net: cli -1077 lines; the crawler + relay machinery are now reusable by the Android app. Adds ConcurrentCollectionsTest; verified via JVM + commonMain metadata compile, the wot/graperank suites, and a bounded live crawl. --- .../com/vitorpamplona/amethyst/cli/Context.kt | 217 +---- .../amethyst/cli/commands/GrapeRankCommand.kt | 639 +------------- .../graperank/GrapeRankDataCrawler.kt | 813 ++++++++++++++++++ .../accessories}/AdaptiveRelayLimiter.kt | 72 +- .../relay/client/accessories/DrainFailure.kt | 77 ++ .../quartz/utils/concurrent/ConcurrentMap.kt | 68 ++ .../quartz/utils/concurrent/ConcurrentSet.kt | 43 + .../concurrent/ConcurrentCollectionsTest.kt | 108 +++ .../concurrent/ConcurrentMap.jvmAndroid.kt | 51 ++ .../concurrent/ConcurrentSet.jvmAndroid.kt | 35 + .../utils/concurrent/ConcurrentMap.native.kt | 80 ++ .../utils/concurrent/ConcurrentSet.native.kt | 46 + 12 files changed, 1406 insertions(+), 843 deletions(-) create mode 100644 quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankDataCrawler.kt rename {cli/src/main/kotlin/com/vitorpamplona/amethyst/cli => quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories}/AdaptiveRelayLimiter.kt (80%) create mode 100644 quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/DrainFailure.kt create mode 100644 quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.kt create mode 100644 quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentSet.kt create mode 100644 quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentCollectionsTest.kt create mode 100644 quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.jvmAndroid.kt create mode 100644 quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentSet.jvmAndroid.kt create mode 100644 quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.native.kt create mode 100644 quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentSet.native.kt 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 4214dac422..531fcbf6c5 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt @@ -41,6 +41,9 @@ import com.vitorpamplona.quartz.nip01Core.core.hexToByteArray import com.vitorpamplona.quartz.nip01Core.crypto.verify import com.vitorpamplona.quartz.nip01Core.metadata.MetadataEvent 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.fetchAllPagesFromPool import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.publishAndConfirmDetailed import com.vitorpamplona.quartz.nip01Core.relay.client.auth.RelayAuthenticator @@ -74,86 +77,13 @@ import kotlinx.coroutines.CompletableDeferred import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.channels.Channel.Factory.UNLIMITED import kotlinx.coroutines.coroutineScope -import kotlinx.coroutines.joinAll import kotlinx.coroutines.launch import kotlinx.coroutines.selects.select import kotlinx.coroutines.withTimeoutOrNull import okhttp3.Dispatcher import okhttp3.OkHttpClient -import java.util.concurrent.ConcurrentHashMap import java.util.concurrent.TimeUnit -/** - * Why a relay could not be used for a drain — when the reason is worth acting on. - * - * - [HARD]: the relay answered wrong, or cannot exist. A bad HTTP upgrade (not a - * websocket / dead status code), an unresolvable domain, or a TLS misconfig. - * This will not fix itself, so one strike is enough to drop it. - * - [TRANSIENT]: a failure that might clear — connection refused / reset, host - * unreachable, or a temporary 429/5xx on the upgrade. Struck a few times - * before we give up. - * - * A pure connect **timeout** is neither. The relay is most likely just busy, so - * we retry it and never mark it dead — [classifyDrainFailure] returns null for - * it (and for any non-failure terminal reason). - */ -enum class DrainFailure { HARD, TRANSIENT } - -/** - * Classify a [Context.drain] per-relay terminal reason. Returns null when the - * relay should simply be retried (a timeout, or a non-failure like eose/closed). - * The reason shape is `cannot:` for a connect failure (see - * `BasicRelayClient.onCannotConnect`), or `eose` / `closed:…` / `timeout`. - */ -fun classifyDrainFailure(reason: String): DrainFailure? { - if (!reason.startsWith("cannot")) return null - val m = reason.removePrefix("cannot:").lowercase() - // The message now carries the exception class name (see BasicRelayClient), so - // we can key on the stable *type* rather than localized message text. - // Busy, not dead: a connect/read timeout means the handshake just didn't - // finish in time. Retry it — the relay is probably fine, only slow or loaded. - if ("timeout" in m || "timed out" in m) return null // SocketTimeoutException, etc. - // Cannot ever work: unresolvable domain (DNS) or a TLS misconfiguration. - // Dead for good — one strike is enough. - if ("unknownhost" in m || // UnknownHostException - "unable to resolve host" in m || - "no address associated" in m || - "nodename nor servname" in m || - "sslhandshake" in m || // SSLHandshakeException - "sslpeerunverified" in m || - "sslexception" in m || - "certificate" in m || // CertificateException - "trust anchor" in m || - "certpath" in m - ) { - return DrainFailure.HARD - } - // Wrong HTTP upgrade. Usually a misconfigured endpoint (not a relay), but - // 429 / 5xx mean "busy, come back later", so those stay transient. - if ("server misconfigured" in m || "not a websocket" in m || "expected http 101" in m) { - val transientCode = Regex("response: (429|500|502|503|504)").containsMatchIn(m) - return if (transientCode) DrainFailure.TRANSIENT else DrainFailure.HARD - } - // Refused / reset / unreachable / anything else: might clear — retry a few times. - return DrainFailure.TRANSIENT -} - -/** - * Max total "entries" (authors + ids + tag values) allowed in a single REQ frame. - * A REQ carries all of a subscription's filters at once, and each entry is a - * ~67-byte JSON hex string, so 2500 entries ≈ 167KB — comfortably under the 256KB - * message cap most relays enforce (and reject a frame over, dropping every author - * in it). [Context.drainGated] groups a relay's filters to stay within this. - */ -private const val MAX_REQ_ENTRIES = 2500 - -/** Count the size-driving entries in a filter: authors, ids, and tag values. */ -fun filterEntries(f: Filter): Int = - (f.authors?.size ?: 0) + - (f.ids?.size ?: 0) + - (f.tags?.values?.sumOf { it.size } ?: 0) + - (f.tagsAll?.values?.sumOf { it.size } ?: 0) - /** * Per-invocation wiring. Each CLI run constructs a Context, does its work, * and then closes it — no daemon. @@ -550,10 +480,8 @@ class Context( timeoutMs: Long = 8_000, diagnoseSlow: Boolean = false, deadOut: MutableMap? = null, - gatePerRelay: Boolean = false, ): List> { if (filters.isEmpty()) return emptyList() - if (gatePerRelay) return drainGated(filters, timeoutMs, diagnoseSlow, deadOut) val eventChannel = Channel>(UNLIMITED) // Carries the terminal reason per relay so a timeout can distinguish a slow // relay (never terminal) from a connect failure / CLOSED. @@ -636,145 +564,6 @@ class Context( return collected } - /** - * Per-relay-gated variant of [drain] used by the crawl. Instead of one - * subscription spanning every relay, each relay gets its own subscription - * held behind [relayLimiter], so we never exceed the relay's adaptive - * concurrent-subscription cap. A relay whose cap is full simply waits for one - * of our other subscriptions on it to finish before its REQ goes out; relays - * we haven't upset run at the full starting cap and never wait. - * - * Semantics match [drain] otherwise: verify+store on a single consumer - * (so store writes stay serialized), return events tagged by relay, and - * report hard connect failures into [deadOut]. - */ - private suspend fun drainGated( - filters: Map>, - timeoutMs: Long, - diagnoseSlow: Boolean, - deadOut: MutableMap?, - ): List> { - val eventChannel = Channel>(UNLIMITED) - - // Split each relay's filters into REQ-sized groups. A REQ frame carries ALL - // its filters at once, so a popular relay routed thousands of authors would - // otherwise produce a multi-MB frame that most relays reject outright - // ("message too large") — silently dropping every author in it. Grouping by - // total entry count keeps each REQ well under the common 256KB cap; a relay - // with more authors just gets several smaller REQs, each its own gated sub. - val units = ArrayList>>() - for ((relay, relayFilters) in filters) { - var group = ArrayList() - var entries = 0 - for (f in relayFilters) { - val fe = filterEntries(f) - if (group.isNotEmpty() && entries + fe > MAX_REQ_ENTRIES) { - units.add(relay to group) - group = ArrayList() - entries = 0 - } - group.add(f) - entries += fe - } - if (group.isNotEmpty()) units.add(relay to group) - } - - // Per-relay failure classification, HARD winning over TRANSIENT across a - // relay's several REQ-groups; plus which relays stalled to a timeout. - val failures = ConcurrentHashMap() - val timedOut = ConcurrentHashMap.newKeySet() - - val collected = mutableListOf>() - coroutineScope { - // Single consumer: verify+store serially. One writer, so SeenIds' - // single-writer contract holds. The outbox model (and the wide relay- - // list broadcast) delivers the SAME event from many relays at once; skip - // a duplicate BEFORE the expensive Schnorr verify+store. An id is marked - // seen only after it verifies, so a forged copy (valid id, bad signature) - // delivered first can't suppress the genuine one that follows. - val consumer = - launch { - val seen = SeenIds(initialSlotsPow2 = 12) - for ((relay, event) in eventChannel) { - if (seen.contains(event.id)) continue - if (verifyAndStore(event)) { - seen.add(event.id) - collected.add(relay to event) - } - } - } - // One gated subscription per (relay, REQ-group). The permit is held for - // the group's whole life, so concurrent subs on a relay never exceed its - // adaptive cap. Each group carries its own subId, listener, and terminal - // signal (relay + subId together identify a group, but a per-group - // listener is simplest). - units - .map { (relay, groupFilters) -> - launch { - relayLimiter.withPermit(relay) { - val subId = newSubId() - val done = CompletableDeferred() - val groupListener = - object : SubscriptionListener { - override fun onEvent( - event: Event, - isLive: Boolean, - r: NormalizedRelayUrl, - forFilters: List?, - ) { - eventChannel.trySend(r to event) - } - - override fun onEose( - r: NormalizedRelayUrl, - forFilters: List?, - ) { - done.complete("eose") - } - - override fun onClosed( - message: String, - r: NormalizedRelayUrl, - forFilters: List?, - ) { - done.complete("closed:$message") - } - - override fun onCannotConnect( - r: NormalizedRelayUrl, - message: String, - forFilters: List?, - ) { - done.complete("cannot:$message") - } - } - client.subscribe(subId, mapOf(relay to groupFilters), groupListener) - try { - val reason = withTimeoutOrNull(timeoutMs) { done.await() } ?: "timeout" - if (reason == "timeout") timedOut.add(relay) - classifyDrainFailure(reason)?.let { kind -> - failures.merge(relay, kind) { a, b -> - if (a == DrainFailure.HARD || b == DrainFailure.HARD) DrainFailure.HARD else DrainFailure.TRANSIENT - } - } - } finally { - client.unsubscribe(subId) - } - } - } - }.joinAll() - // All subscriptions are torn down; no more events can arrive. Close the - // channel so the consumer drains what's buffered and completes. - eventChannel.close() - consumer.join() - } - if (diagnoseSlow && timedOut.isNotEmpty()) { - logSlowDrain(timeoutMs, timedOut, emptyMap(), collected) - } - deadOut?.putAll(failures) - return collected - } - /** * On a [drain] timeout, report which relays stalled and why — a relay that * never sent EOSE (slow, possibly still streaming) vs one that couldn't be diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/GrapeRankCommand.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/GrapeRankCommand.kt index 955948b741..39693a5ade 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/GrapeRankCommand.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/GrapeRankCommand.kt @@ -23,11 +23,11 @@ package com.vitorpamplona.amethyst.cli.commands import com.vitorpamplona.amethyst.cli.Args import com.vitorpamplona.amethyst.cli.Context import com.vitorpamplona.amethyst.cli.DataDir -import com.vitorpamplona.amethyst.cli.DrainFailure import com.vitorpamplona.amethyst.cli.Output import com.vitorpamplona.amethyst.commons.defaults.Constants import com.vitorpamplona.amethyst.commons.defaults.DefaultIndexerRelayList import com.vitorpamplona.quartz.experimental.graperank.GrapeRank +import com.vitorpamplona.quartz.experimental.graperank.GrapeRankDataCrawler import com.vitorpamplona.quartz.experimental.graperank.GrapeRankParams import com.vitorpamplona.quartz.experimental.graperank.TrustGraphBuilder import com.vitorpamplona.quartz.nip01Core.core.Event @@ -44,7 +44,6 @@ import com.vitorpamplona.quartz.nip09Deletions.DeletionEvent import com.vitorpamplona.quartz.nip09Deletions.DeletionIndex import com.vitorpamplona.quartz.nip51Lists.muteList.MuteListEvent import com.vitorpamplona.quartz.nip56Reports.ReportEvent -import com.vitorpamplona.quartz.nip65RelayList.AdvertisedRelayListEvent import com.vitorpamplona.quartz.nip85TrustedAssertions.list.TrustProviderListEvent import com.vitorpamplona.quartz.nip85TrustedAssertions.list.serviceProviders import com.vitorpamplona.quartz.nip85TrustedAssertions.list.tags.ProviderTypes @@ -52,18 +51,10 @@ import com.vitorpamplona.quartz.nip85TrustedAssertions.list.tags.ServiceProvider import com.vitorpamplona.quartz.nip85TrustedAssertions.list.tags.ServiceType import com.vitorpamplona.quartz.nip85TrustedAssertions.users.ContactCardEvent import com.vitorpamplona.quartz.nip85TrustedAssertions.users.tags.RankTag -import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers -import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.async import kotlinx.coroutines.awaitAll -import kotlinx.coroutines.cancel -import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.coroutineScope -import kotlinx.coroutines.joinAll -import kotlinx.coroutines.launch -import java.util.concurrent.ConcurrentHashMap -import kotlin.coroutines.coroutineContext import kotlin.math.roundToInt /** @@ -91,9 +82,6 @@ import kotlin.math.roundToInt * - `amy graperank providers [USER]` — list a user's trusted providers. */ object GrapeRankCommand { - // Authors per REQ filter — keeps individual subscriptions within relay limits. - private const val AUTHORS_PER_FILTER = 300 - // Concurrent publishes when writing NIP-85 cards. private const val PUBLISH_CONCURRENCY = 16 @@ -103,54 +91,8 @@ object GrapeRankCommand { // message cap). private const val DELETE_PER_EVENT = 400 - // Times we re-query an unreachable user's outbox before giving up on it, so - // the crawl still terminates on a finite graph. - private const val MAX_OUTBOX_ATTEMPTS = 3 - - // Users whose outboxes we fetch in a single drain. Draining thousands of - // distinct outbox relays at once saturates connections and times out - // (empirically ~250 users/drain succeeds, ~17k fails); keep the fan-out small. - private const val USER_BATCH = 256 - - // Global content-drain fan-out — how many outbox batches we drain at once. - // This is a GLOBAL bound (memory / open sockets); the per-relay concurrency - // limit is enforced separately and adaptively by [AdaptiveRelayLimiter] - // (drains run with gatePerRelay=true), which starts every relay at 100 - // concurrent subs and demotes only the ones that complain (100 → 20 → 10). - // The two compose: at fan-out 24 a well-behaved relay runs at up to 24 - // concurrent subs, while a relay that pushes back is cut to 20 then 10 — - // below the global bound, so the ladder actually bites. A higher global - // fan-out (measured at 48) *re-floods* the busy hubs faster than demotion - // catches up ("max concurrent subscription count reached" spikes) and - // regressed wall-time, so keep the global bound moderate and let the - // per-relay cap do the targeting. - private const val DRAIN_CONCURRENCY = 24 - - // Sharded backbone sweep: instead of asking every popular relay for the - // same full author list (N× redundant), split the still-missing authors - // into SHARD_RELAYS lists and send each to ONE of the top relays. Authors a - // relay doesn't have rotate onto a different relay next pass, up to - // SHARD_ROTATIONS times, so over a few passes each author is tried on - // several popular relays. Once the remaining set drops below - // SHARD_BROADCAST_THRESHOLD it's cheap to just ask them all at once. - private const val SHARD_RELAYS = 10 - private const val SHARD_ROTATIONS = 6 - private const val SHARD_BROADCAST_THRESHOLD = 2000 - - // The small-remainder broadcast (once a sweep is under the threshold) goes to - // this many top live relays, not just the SHARD_RELAYS the rotation used — - // a user's kind:3 is often mirrored on a busy relay ranked below the top 10, - // which is where the old last-mile pass found its stragglers. - private const val BROADCAST_RELAYS = 60 - - // A relay that fails to CONNECT this many times is treated as dead and - // dropped from routing, so we stop paying the drain timeout on it. Kept above - // 1 so a single transient connect blip doesn't evict a relay for the run. - private const val MAX_DEAD_STRIKES = 3 - // Broad, big general relays that carry kind:10002 for many users, added to the - // discovery set to raise the odds of resolving a stranger's outbox. Every entry - // is NIP-11 liveness-checked — dead relays only add timeouts. + // crawler's discovery set to raise the odds of resolving a stranger's outbox. private val EXTRA_DISCOVERY_RELAYS: Set = listOf( "wss://relay.damus.io", @@ -160,20 +102,6 @@ object GrapeRankCommand { "wss://eden.nostr.land", ).mapNotNull { RelayUrlNormalizer.normalizeOrNull(it) }.toSet() - // How many of the most-used write relays (learned from everyone's kind:10002) - // to keep as the known-good backbone for retrying users we couldn't reach. - private const val BACKBONE_SIZE = 30 - - // Warm pool: hold a persistent, do-nothing subscription open to the busiest - // WARM_POOL_SIZE relays for the whole crawl, so the connections we reuse - // every round survive the between-round routing gaps (and niche-relay churn) - // instead of being dropped ~300ms after a wave ends and reconnected next - // round. The filter matches an impossible event id, so the relay EOSEs - // immediately and streams nothing — it only keeps the socket warm. - private const val WARM_POOL_SIZE = 20 - private const val WARM_SUB_ID = "graperank-warm" - private val WARM_FILTERS = listOf(Filter(ids = listOf("0".repeat(64)))) - suspend fun dispatch( dataDir: DataDir, tail: Array, @@ -225,354 +153,47 @@ object GrapeRankCommand { ctx.prepare() val observer = observerArg?.let { ctx.requireUserHex(it) } ?: ctx.identity.pubKeyHex - val graphKinds = listOf(ContactListEvent.KIND, MuteListEvent.KIND, ReportEvent.KIND) - // Kinds requested from relays during the crawl: the graph edges PLUS the - // user's own kind:10002. A user's outbox holds the freshest copy of their - // relay list, so folding 10002 into the same query we send their outbox - // keeps our routing current instead of trusting a possibly-stale indexer - // copy. Safe to also pull from popular relays in the sweep — the store - // keeps newest-by-created_at for the replaceable 10002, so the freshest - // always wins regardless of which relay delivered it. - val fetchKinds = graphKinds + AdvertisedRelayListEvent.KIND - - // The graph is built incrementally: contact lists stream straight into a - // compact int-CSR structure and the Event is discarded, so the whole - // network fits in memory without holding millions of kind:3 objects. + // Contact lists stream straight into a compact int-CSR structure as the + // crawl finds them and the Event is discarded, so the whole network fits + // in memory without holding millions of kind:3 objects. val builder = TrustGraphBuilder() - var rounds = 0 - var relaysContactedCount = 0 var contactListsFed = 0 // Wall time to read + deserialize the contact lists out of the store - // (offline path only; online streams them in during the crawl). This - // is the real pre-scoring cost — the int-CSR build afterwards is a - // cheap in-memory pack. + // (offline path only; online streams them in during the crawl). var storeLoadMs: Long? = null - // Wall time to crawl + download the whole graph off the relays - // (online path only) — rounds + last-mile sweep, i.e. everything up - // to the point the graph is fully fetched. This is network-bound and - // dominates a from-scratch run. - var downloadMs: Long? = null + // Crawl telemetry (online path only): rounds, relays contacted, the + // per-hop histogram, and the network-bound download time that dominates a + // from-scratch run. Null on the offline path. + var crawlStats: GrapeRankDataCrawler.Stats? = null - val hopOf = HashMap() if (!offline) { - val crawlStart = System.nanoTime() - // Scope for fire-and-forget relay-list discovery: the wide Tier-2 - // sweep (ensureRelayLists) casts kind:10002 queries across every relay - // we know, but we don't block the crawl on it — its results just - // enrich routing for later rounds. SupervisorJob so one failing sweep - // never cancels the others; cancelled when the crawl finishes. - val bgScope = CoroutineScope(coroutineContext + SupervisorJob()) - val discovered = hashSetOf(observer) - hopOf[observer] = 0 - // Per-user relay hints harvested from the `p`-tag relay hints in the - // contact lists we crawl (A's follow of B says where B writes) — a - // discovery tier below each user's kind:10002 outbox. Concurrent: - // the Phase-B producer reads these while the consumer's ingest writes - // them (see the worker-pool below), so both map and inner sets are - // thread-safe. - val relayHints = ConcurrentHashMap>() - // Users we're finished with this run: we fed their latest kind:3, or - // ran out of retry attempts on an unreachable outbox. - val done = hashSetOf() - // Outbox retry counts. Concurrent: the producer reads (to widen a - // retry's routing) while the consumer increments. - val attempts = ConcurrentHashMap() - val relaysContacted = hashSetOf() - // Known-good relay pool, learned from the crawl itself: how often each - // relay appears as someone's write relay, and which relays actually - // delivered events (so we know they connect and work). The most-common - // live relays become the `backbone` we retry unreachable users against. - val writeRelayFreq = HashMap() - val liveRelays = hashSetOf() - // Relays that failed to connect MAX_DEAD_STRIKES times — dropped - // from all routing so a wave stops eating the timeout on them. - // Concurrent: drain workers strike relays while the producer reads - // deadRelays to prune routing. - val deadRelays = ConcurrentHashMap.newKeySet() - val relayStrikes = ConcurrentHashMap() - - // A relay that HARD-failed (bad domain, TLS misconfig, dead HTTP - // code — see DrainFailure) is dropped on the first strike: it will - // not fix itself. A TRANSIENT failure (refused/reset/unreachable, - // or a 429/5xx) might clear, so it takes MAX_DEAD_STRIKES before we - // give up. Pure timeouts never reach here — the drain treats them as - // busy-retry and does not report them dead at all. - fun recordDead(failed: Map) { - for ((r, kind) in failed) { - when (kind) { - DrainFailure.HARD -> deadRelays.add(r) - DrainFailure.TRANSIENT -> - if (relayStrikes.merge(r, 1, Int::plus)!! >= MAX_DEAD_STRIKES) deadRelays.add(r) - } - } - } - - // The busiest live relays we've learned, excluding the dead ones. - fun topLiveRelays(cap: Int): List = - writeRelayFreq.entries - .asSequence() - .filter { it.key in liveRelays && it.key !in deadRelays } - .sortedByDescending { it.value } - .take(cap) - .map { it.key } - .toList() - - // Feed a user's contact list into the graph, harvest relay hints, stamp - // the hop distance of newly-seen follows, and add them to the frontier. - // Called once per user (guarded by `done`). Returns the count of - // newly-discovered users. - fun ingest( - source: HexKey, - contacts: ContactListEvent, - ): Int { - val nextHop = (hopOf[source] ?: 0) + 1 - val follows = ArrayList() - var fresh = 0 - for (tag in contacts.follows()) { - follows.add(tag.pubKey) - tag.relayUri?.let { relayHints.getOrPut(tag.pubKey) { ConcurrentHashMap.newKeySet() }.add(it) } - if (discovered.add(tag.pubKey)) { - hopOf[tag.pubKey] = nextHop - fresh++ - } - } - builder.addFollows(source, follows) - contactListsFed++ - return fresh - } - - // Feed into the graph the contact lists a drain just returned - // (deduped by author; the store's canonical latest wins), marking - // fed authors done. Only the authors we actually received are - // touched — no scan over the whole still-missing set. Returns the - // count newly fed. - suspend fun harvest(events: List>): Int { - var got = 0 - for ((_, ev) in events) { - if (ev !is ContactListEvent) continue - val pk = ev.pubKey - if (pk in done) continue - val contacts = ctx.contactsOf(pk) ?: continue - done += pk - ingest(pk, contacts) - got++ - } - return got - } - - // Sharded backbone sweep (see SHARD_RELAYS). Splits the missing - // authors across the top live relays — one shard per relay, so no - // relay gets the same list twice — drains all shards concurrently, - // then rotates whoever's still missing onto a different relay for up - // to SHARD_ROTATIONS passes. Once the remainder is small it's cheap - // to broadcast it to every top relay at once. Returns lists fed. - suspend fun shardedSweep(authors: Collection): Int { - val top = topLiveRelays(SHARD_RELAYS) - if (top.isEmpty()) return 0 - val n = top.size - var missing = authors.filter { it !in done && ctx.contactsOf(it) == null } - var got = 0 - var rotation = 0 - while (missing.size > SHARD_BROADCAST_THRESHOLD && rotation < SHARD_ROTATIONS) { - val shards = Array(n) { ArrayList() } - for (pk in missing) { - val base = ((pk.hashCode() % n) + n) % n - shards[(base + rotation) % n].add(pk) - } - val results = - coroutineScope { - top - .mapIndexedNotNull { i, relay -> - val shard = shards[i] - if (shard.isEmpty()) { - null - } else { - // Each drain gets its own dead-set — the concurrent - // drains must not share a mutable HashSet. - async { - val dead = HashMap() - val filters = - mapOf(relay to shard.chunked(AUTHORS_PER_FILTER).map { Filter(kinds = fetchKinds, authors = it) }) - ctx.drain(filters, timeoutMs, diagnose, dead, gatePerRelay = true) to dead - } - } - }.awaitAll() - } - for ((_, dead) in results) recordDead(dead) - relaysContacted += top - val flat = results.flatMap { it.first } - for ((relay, _) in flat) liveRelays.add(relay) - got += harvest(flat) - missing = missing.filter { it !in done } - rotation++ - } - // Once the remainder is small it's cheap to ask every top relay - // for it at once. If the rotations bailed with a still-large set, - // those authors just aren't on the popular relays — leave them to - // the caller's outbox pass rather than broadcast a huge list. - if (missing.isNotEmpty() && missing.size <= SHARD_BROADCAST_THRESHOLD) { - // Broadcast the small remainder to a wider set of busy relays - // than the rotation used — recovers users whose list is only - // on a relay ranked below the top SHARD_RELAYS. - val live = topLiveRelays(BROADCAST_RELAYS) - if (live.isNotEmpty()) { - val dead = HashMap() - val filters = - live.associateWith { missing.chunked(AUTHORS_PER_FILTER).map { Filter(kinds = fetchKinds, authors = it) } } - val events = ctx.drain(filters, timeoutMs, diagnose, dead, gatePerRelay = true) - recordDead(dead) - relaysContacted += live - for ((relay, _) in events) liveRelays.add(relay) - got += harvest(events) - } - } - return got - } - - // Crawl to full graph depth (no user cap; --max-hops bounds the follow - // distance). Each run fetches every discovered user's LATEST - // kind:3/10000/1984 once from their outbox (a freshness pass — grouped - // by write relay in routeByOutbox), unless we already fetched it this - // run (`done`). An unreachable outbox is retried up to - // MAX_OUTBOX_ATTEMPTS then dropped so the crawl terminates. - while (rounds < maxRounds) { - // Only crawl users within the hop budget; deeper users still appear - // in the graph as follow targets, we just don't fetch their lists. - val pending = discovered.filter { it !in done && (hopOf[it] ?: 0) < maxHops } - if (pending.isEmpty()) break - rounds++ - - // Refresh the warm pool to this round's busiest relays and keep - // that subscription open — reusing the same subId just updates the - // desired-relay set, so these sockets stay up across the round. - topLiveRelays(WARM_POOL_SIZE).takeIf { it.isNotEmpty() }?.let { warm -> - ctx.client.subscribe(WARM_SUB_ID, warm.associateWith { WARM_FILTERS }, null) - } - - val discoveredBefore = discovered.size - val fedBefore = contactListsFed - - // Phase A — bulk-fetch from the busiest relays via the sharded - // sweep. Most users' kind:3 lives on the big popular relays, so - // this clears the majority cheaply, without asking every relay for - // the same authors (early rounds no-op until a backbone is learned). - shardedSweep(pending) - - // Phase B — whoever the popular relays didn't have (niche - // outboxes): resolve their kind:10002, then fetch from their own - // write relays, drained a few at a time and skipping dead relays. - val stragglers = pending.filter { it !in done } - if (stragglers.isNotEmpty()) { - val backbone = topLiveRelays(BACKBONE_SIZE).toSet() - // Snapshot of every relay we've seen work, for the wide Tier-2 - // sweep (taken now, on this single coroutine, before the Phase-B - // workers start mutating liveRelays). - val allLive = (liveRelays - deadRelays).toSet() - ensureRelayLists(ctx, stragglers.toSet(), allLive, bgScope, timeoutMs, diagnose) - - // Continuous worker pool instead of chunked awaitAll barriers. - // The old shape drained DRAIN_CONCURRENCY batches, waited for the - // SLOWEST (a dead relay's full timeout), ingested, then started - // the next group — so every batch's tail idled the whole pool. - // Here a fixed set of DRAIN_CONCURRENCY workers pulls batches off - // a queue and grabs the next the instant a drain returns, so no - // worker waits on a slow sibling and hot relays stay connected - // (some worker is always subscribed). Shared graph state stays - // single-writer: routeByOutbox runs only on the producer (keeps - // writeRelayFreq serial) and ingest runs only on the consumer - // (keeps discovered/done/builder/hopOf serial), now overlapped - // with draining instead of blocked behind each batch. - val routed = Channel, Map>>>(DRAIN_CONCURRENCY * 2) - val drainedOut = Channel, Set, List>>>(Channel.UNLIMITED) - coroutineScope { - // Producer: route each batch by outbox (serial), backpressured - // by the bounded `routed` channel so we don't precompute every - // filter map at once. - val producer = - launch { - for (batch in stragglers.chunked(USER_BATCH)) { - val filters = routeByOutbox(ctx, batch.toSet(), relayHints, backbone, attempts, writeRelayFreq, fetchKinds, deadRelays) - routed.send(batch to filters) - } - routed.close() - } - // Drain workers: pure network, no shared graph-state writes - // except recordDead (concurrent-safe now). - val workers = - List(DRAIN_CONCURRENCY) { - launch { - for ((batch, filters) in routed) { - val dead = HashMap() - val events = ctx.drain(filters, timeoutMs, diagnose, dead, gatePerRelay = true) - recordDead(dead) - drainedOut.send(Triple(batch, filters.keys, events)) - } - } - } - // Consumer: single-writer ingest, overlapped with draining. - val consumer = - launch { - for ((batch, relays, events) in drainedOut) { - relaysContacted += relays - // Any relay that gave us an event is proven live + useful. - for ((relay, _) in events) liveRelays.add(relay) - for (pk in batch) { - if (pk in done) continue - val contacts = ctx.contactsOf(pk) - if (contacts != null) { - done += pk - ingest(pk, contacts) - } else { - val tries = (attempts[pk] ?: 0) + 1 - attempts[pk] = tries - if (tries >= MAX_OUTBOX_ATTEMPTS) done += pk - } - } - } - } - producer.join() - workers.joinAll() - drainedOut.close() - consumer.join() - } - } - - System.err.println( - "[graperank] round $rounds: pending=${pending.size}, " + - "gotList=${contactListsFed - fedBefore}, newUsers=${discovered.size - discoveredBefore}, " + - "discovered=${discovered.size}, done=${done.size}, dead=${deadRelays.size}", + // Relay policy for the crawler — where a stranger's kind:10002 is + // found (index/discovery aggregators + general defaults that carry + // kind:10002 for most of the network) and the best-effort general + // relays that might hold content when an outbox is unknown. These + // defaults live in app code, so the quartz crawler takes them injected. + val discoveryRelays = + ctx.bootstrapRelays() + Constants.eventFinderRelays + DefaultIndexerRelayList + EXTRA_DISCOVERY_RELAYS + val contentFallback = ctx.bootstrapRelays() + Constants.eventFinderRelays + val crawler = + GrapeRankDataCrawler( + client = ctx.client, + store = ctx.store, + limiter = ctx.relayLimiter, + config = + GrapeRankDataCrawler.Config( + relayListDiscoveryRelays = discoveryRelays, + contentFallbackRelays = contentFallback, + maxRounds = maxRounds, + maxHops = maxHops, + timeoutMs = timeoutMs, + diagnose = diagnose, + ), + log = { System.err.println(it) }, ) - } - - // Crawl done — drop the warm pool and stop any background relay-list - // sweeps still in flight (their results are already in the store). - ctx.client.unsubscribe(WARM_SUB_ID) - bgScope.cancel() - - // Reports can be retracted. Ask each reporter's outbox for NIP-09 - // kind:5 deletions that cite the reports we gathered (#e-filtered to - // our report ids — not every deletion the user ever made). A report - // the author has since deleted must not count as a negative edge; - // [materializeReports] drops those below. - fetchReportDeletions(ctx, topLiveRelays(BACKBONE_SIZE).toSet(), deadRelays, timeoutMs, diagnose) - - // No separate last-mile pass: the per-round sharded sweep already - // broadcasts the small remaining set to every top relay once it drops - // below SHARD_BROADCAST_THRESHOLD, and the round loop only exits when - // every reachable user within the hop budget is done. - - relaysContactedCount = relaysContacted.size - val perHop = - hopOf.values - .groupingBy { it } - .eachCount() - .toSortedMap() - downloadMs = (System.nanoTime() - crawlStart) / 1_000_000 - System.err.println( - "[graperank] crawl complete: ${discovered.size} discovered, $contactListsFed contact lists fed, " + - "$relaysContactedCount relays contacted, ${deadRelays.size} dead, $rounds rounds in $downloadMs ms; " + - "by hop: " + perHop.entries.joinToString(" ") { "${it.key}=${it.value}" }, - ) + val stats = crawler.crawl(observer, builder) + crawlStats = stats + contactListsFed = stats.contactListsFed if (ctx.relayDiagnostics.hadFeedback()) { System.err.println("[graperank] relay feedback: ${ctx.relayDiagnostics.snapshot()}") } @@ -631,22 +252,17 @@ object GrapeRankCommand { val result = linkedMapOf( "observer" to observer, - "crawl_rounds" to rounds, - "relays_contacted" to relaysContactedCount, + "crawl_rounds" to (crawlStats?.rounds ?: 0), + "relays_contacted" to (crawlStats?.relaysContacted ?: 0), "relay_feedback" to if (ctx.relayDiagnostics.hadFeedback()) ctx.relayDiagnostics.snapshot() else null, "relay_throttling" to if (ctx.relayLimiter.hadThrottling()) ctx.relayLimiter.snapshot() else null, - "max_hop_reached" to (hopOf.values.maxOrNull() ?: 0), - "users_by_hop" to - hopOf.values - .groupingBy { it } - .eachCount() - .toSortedMap() - .mapKeys { it.key.toString() }, + "max_hop_reached" to (crawlStats?.hopHistogram?.keys?.maxOrNull() ?: 0), + "users_by_hop" to (crawlStats?.hopHistogram?.mapKeys { it.key.toString() } ?: emptyMap()), "graph_users" to graph.nodeCount, "graph_edges" to graph.edgeCount(), "reports_deleted" to reportsDeleted, "users_scored" to rankedIds.size, - "download_ms" to downloadMs, + "download_ms" to crawlStats?.downloadMs, "store_load_ms" to storeLoadMs, "graph_build_ms" to buildMs, "scoring_ms" to scoringMs, @@ -1010,133 +626,6 @@ object GrapeRankCommand { return providerListOf(ctx, pubKey) } - /** - * Relays to query for **kind:10002 relay lists** — the account's own relays + - * bootstrap defaults + event-finder relays + the **indexer relays** - * (purplepag.es, coracle, …). Indexers aggregate kind:10002 (and kind:0) for - * the whole network, so this is where a stranger's relay list is found. They - * do NOT hold kind:3/10000/1984 — see [contentFallbackRelays]. - */ - private suspend fun relayListDiscoveryRelays(ctx: Context): Set = ctx.bootstrapRelays() + Constants.eventFinderRelays + DefaultIndexerRelayList + EXTRA_DISCOVERY_RELAYS - - /** - * Best-effort fallback relays for **content** (kind:3/10000/1984/0) when a - * user's outbox is unknown or unreachable. Content lives on each user's own - * outbox, so this is only general-purpose relays that *might* hold a copy — - * bootstrap + event-finder. **No indexers**: they don't serve these kinds. - */ - private suspend fun contentFallbackRelays(ctx: Context): Set = ctx.bootstrapRelays() + Constants.eventFinderRelays - - /** - * Fetch kind:10002 relay lists for any [pubkeys] we don't already know, so - * [routeByOutbox] can route their content query to their own write relays. - * - * Tier 1 queries the bounded relay-list discovery set (indexers + general - * defaults), which aggregate kind:10002 for the whole network — reliable in - * bulk, unlike fanning out to thousands of per-user outboxes. - * - * Tier 2 is a completeness net for the stragglers the indexers don't cover: - * a user publishes their own kind:10002 to their own write relays, and those - * relays overlap heavily with [fallbackRelays] — the known-good backbone we - * learned from the `r` tags in *everyone else's* 10002s. So after tier 1, - * any pubkey still without a relay list is retried against that learned pool - * (minus the tier-1 relays we already asked). Early rounds skip tier 2 - * harmlessly because the backbone is still empty; it kicks in once the crawl - * has learned which relays actually carry 10002s. - */ - private suspend fun ensureRelayLists( - ctx: Context, - pubkeys: Set, - allLiveRelays: Set, - bgScope: CoroutineScope, - timeoutMs: Long, - diagnose: Boolean, - ) { - val missing = pubkeys.filter { ctx.relaysOf(it) == null } - if (missing.isEmpty()) return - - suspend fun query( - authors: List, - relays: Set, - ) { - if (relays.isEmpty() || authors.isEmpty()) return - val filters = - relays.associateWith { - authors.chunked(AUTHORS_PER_FILTER).map { chunk -> - Filter(kinds = listOf(AdvertisedRelayListEvent.KIND), authors = chunk) - } - } - ctx.drain(filters, timeoutMs, diagnose, gatePerRelay = true) - } - - // Tier 1: the index/discovery aggregators, which carry kind:10002 for most - // of the network. Blocking, because this round's routing needs the result. - val discovery = relayListDiscoveryRelays(ctx) - query(missing, discovery) - - // Tier 2: whoever the aggregators still don't have, cast the widest net — - // ask EVERY relay we've seen deliver events, not just the backbone. Fired - // fire-and-forget on [bgScope]: a stray 10002 might sit on any one relay, so - // we don't want to skip any, but we also can't block the crawl on a fan-out - // that large. The results land in the store and improve routing for later - // rounds; anyone still unresolved is handled by fallback routing meanwhile. - val stillMissing = missing.filter { ctx.relaysOf(it) == null } - val wide = allLiveRelays - discovery - if (stillMissing.isNotEmpty() && wide.isNotEmpty()) { - bgScope.launch { query(stillMissing, wide) } - } - } - - /** - * Fetch NIP-09 kind:5 deletion requests that retract any report we gathered. - * - * A reporter can delete their own kind:1984 report. That deletion is valid - * only if it comes from the reporter's own key, and it's published to the - * reporter's outbox — so we group report ids by their author and ask each - * author's write relays for kind:5 events that cite those ids (`#e`). That - * `#e` filter is the point: we pull only the deletions that touch our reports, - * not every deletion the user has ever made. The events land in the store; - * [materializeReports] decides which reports they actually retract. - */ - private suspend fun fetchReportDeletions( - ctx: Context, - backbone: Set, - deadRelays: Set, - timeoutMs: Long, - diagnose: Boolean, - ) { - val idsByAuthor = HashMap>() - for (ev in ctx.store.query(Filter(kinds = listOf(ReportEvent.KIND)))) { - if (ev is ReportEvent) idsByAuthor.getOrPut(ev.pubKey) { ArrayList() }.add(ev.id) - } - if (idsByAuthor.isEmpty()) return - - // Route each reporter to their own write relays (fallback: backbone). - val perRelayAuthors = HashMap>() - for (author in idsByAuthor.keys) { - val write = ctx.relaysOf(author)?.writeRelaysNorm()?.takeIf { it.isNotEmpty() } ?: backbone - for (relay in write) if (relay !in deadRelays) perRelayAuthors.getOrPut(relay) { HashSet() }.add(author) - } - if (perRelayAuthors.isEmpty()) return - - val filters = - perRelayAuthors.mapValues { (_, authors) -> - buildList { - for (authorChunk in authors.chunked(AUTHORS_PER_FILTER)) { - // Scope #e to this author-chunk's own report ids, chunked to - // respect REQ limits. Any over-match (a filter pairing an - // author with another author's id) is harmless — the - // deleter-must-be-author check in materializeReports rejects it. - val chunkIds = authorChunk.flatMap { idsByAuthor[it].orEmpty() } - for (idChunk in chunkIds.chunked(AUTHORS_PER_FILTER)) { - add(Filter(kinds = listOf(DeletionEvent.KIND), authors = authorChunk, tags = mapOf("e" to idChunk))) - } - } - } - } - ctx.drain(filters, timeoutMs, diagnose, gatePerRelay = true) - } - /** * Feed reports into [builder], dropping any that a valid NIP-09 deletion has * retracted. Uses quartz's [DeletionIndex] — the same indexer the Android @@ -1171,52 +660,6 @@ object GrapeRankCommand { return dropped } - /** - * Group [pubkeys] by the relays we should query for their events: - * - first try: the user's own kind:10002 write relays (the outbox model); - * - a retry (`attempts[pk] > 0`, its outbox already failed): outbox + - * [backbone] — the known-good relays other people write to, which likely - * hold a copy; - * - no outbox at all: harvested [hints] + backbone + the general fallback. - * - * Also tallies each user's write relays into [writeRelayFreq] so the backbone - * can be learned from the crawl. Authors are chunked per relay to respect REQ - * limits. - */ - private suspend fun routeByOutbox( - ctx: Context, - pubkeys: Set, - hints: Map>, - backbone: Set, - attempts: Map, - writeRelayFreq: MutableMap, - kinds: List, - deadRelays: Set, - ): Map> { - val fallback = contentFallbackRelays(ctx) - val perRelay = HashMap>() - - for (pk in pubkeys) { - val write = ctx.relaysOf(pk)?.writeRelaysNorm()?.takeIf { it.isNotEmpty() } - write?.forEach { writeRelayFreq.merge(it, 1, Int::plus) } - val relays = - when { - write == null -> hints[pk].orEmpty() + backbone + fallback - (attempts[pk] ?: 0) > 0 -> write + backbone - else -> write - } - // Skip relays already proven dead — routing to them only burns the - // drain timeout. - for (relay in relays) if (relay !in deadRelays) perRelay.getOrPut(relay) { HashSet() }.add(pk) - } - - return perRelay.mapValues { (_, authors) -> - authors.chunked(AUTHORS_PER_FILTER).map { chunk -> - Filter(kinds = kinds, authors = chunk) - } - } - } - /** * The exact `rank` tag VALUE STRING we last published for each target, read * from the active account's own kind:30382 cards in the local store (newest 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 new file mode 100644 index 0000000000..f25c177010 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankDataCrawler.kt @@ -0,0 +1,813 @@ +/* + * 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.experimental.graperank + +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.core.HexKey +import com.vitorpamplona.quartz.nip01Core.crypto.verify +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.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.store.IEventStore +import com.vitorpamplona.quartz.nip02FollowList.ContactListEvent +import com.vitorpamplona.quartz.nip09Deletions.DeletionEvent +import com.vitorpamplona.quartz.nip51Lists.muteList.MuteListEvent +import com.vitorpamplona.quartz.nip56Reports.ReportEvent +import com.vitorpamplona.quartz.nip65RelayList.AdvertisedRelayListEvent +import com.vitorpamplona.quartz.utils.Log +import com.vitorpamplona.quartz.utils.SeenIds +import com.vitorpamplona.quartz.utils.concurrent.ConcurrentMap +import com.vitorpamplona.quartz.utils.concurrent.ConcurrentSet +import kotlinx.coroutines.CompletableDeferred +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.async +import kotlinx.coroutines.awaitAll +import kotlinx.coroutines.cancel +import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.coroutineScope +import kotlinx.coroutines.joinAll +import kotlinx.coroutines.launch +import kotlinx.coroutines.withTimeoutOrNull +import kotlin.coroutines.coroutineContext +import kotlin.time.TimeSource + +/** + * Crawls the Nostr follow/mute/report graph outward from an observer and streams + * the contact lists it finds into a [TrustGraphBuilder], so [GrapeRank] can score + * the whole reachable network from that observer's point of view. + * + * It uses the outbox model: each user's kind:10002 write relays are located + * first, then their kind:3 / kind:10000 / kind:1984 events are fetched from + * *their own* relays. The crawl is exhaustive — no user cap; it keeps going until + * every discovered user's outbox has been checked and their contact list pulled + * (an unreachable outbox is retried a few times), bounded only by [Config.maxHops] + * (follow-graph distance) and the [Config.maxRounds] safety backstop. + * + * Every event it fetches (contact lists, mute lists, reports, relay lists, and + * the report deletions it looks up) is verified and persisted to [store], so the + * caller can materialize mutes + reports (honouring NIP-09 retractions) from the + * store afterwards. Only the contact lists are streamed into the [TrustGraphBuilder] + * during the crawl — the compact int-CSR structure keeps the whole network in + * memory without holding millions of kind:3 objects. + * + * The crawler is transport-agnostic within quartz: it takes a [NostrClient], an + * [IEventStore], and the shared [AdaptiveRelayLimiter] (which must already be + * registered as a connection listener on the client so its ladders react to + * NOTICE/CLOSED frames). Relay *policy* — which aggregators know kind:10002, which + * general relays might hold content — is injected via [Config], because those + * defaults live in application code, not the protocol library. Operator progress + * is emitted through [log]; a headless caller routes it to stderr, a UI ignores it. + */ +class GrapeRankDataCrawler( + private val client: NostrClient, + private val store: IEventStore, + private val limiter: AdaptiveRelayLimiter, + private val config: Config, + private val log: (String) -> Unit = {}, +) { + /** + * Relay policy + crawl bounds. The relay sets come from the caller because the + * aggregator/bootstrap defaults live outside quartz. + * + * @param relayListDiscoveryRelays where to look up a stranger's kind:10002 — + * the index/discovery aggregators (purplepag.es, coracle, …) plus general + * defaults that carry kind:10002 for most of the network. + * @param contentFallbackRelays best-effort general relays that *might* hold a + * user's kind:3/10000/1984 when their outbox is unknown or unreachable. + * @param maxRounds safety backstop on freshness passes (default: run to convergence). + * @param maxHops follow-graph distance from the observer to crawl (Brainstorm uses 8). + * @param timeoutMs per-drain timeout. + * @param diagnose log a breakdown of slow/unreachable relays on each drain timeout. + */ + class Config( + val relayListDiscoveryRelays: Set, + val contentFallbackRelays: Set, + val maxRounds: Int = Int.MAX_VALUE, + val maxHops: Int = Int.MAX_VALUE, + val timeoutMs: Long = 10_000, + val diagnose: Boolean = false, + ) + + /** What the crawl fetched — the counters the caller reports and the graph is built from. */ + class Stats( + val rounds: Int, + val discovered: Int, + val contactListsFed: Int, + val relaysContacted: Int, + val deadRelays: Int, + /** Users bucketed by follow-graph distance from the observer (hop -> count), ascending. */ + val hopHistogram: Map, + val downloadMs: Long, + ) + + /** + * Crawl from [observer], streaming discovered contact lists into [builder] + * (follows only — mutes/reports land in the store for the caller to + * materialize). Returns the crawl [Stats]. + */ + suspend fun crawl( + observer: HexKey, + builder: TrustGraphBuilder, + ): Stats = CrawlRun(observer, builder).run() + + /** + * Holds all per-crawl mutable state. Graph state (discovered/done/hopOf/ + * builder/writeRelayFreq/liveRelays/relaysContacted) is single-writer by + * construction — Phase A and the Phase-B consumer never run concurrently, and + * routeByOutbox (the only Phase-B producer write, to writeRelayFreq) touches a + * disjoint field — so those stay plain collections. Only the state genuinely + * shared across the producer / consumer / drain-worker coroutines is concurrent: + * relayHints, attempts, deadRelays, relayStrikes. + */ + private inner class CrawlRun( + val observer: HexKey, + val builder: TrustGraphBuilder, + ) { + val hopOf = HashMap() + val discovered = hashSetOf(observer) + val done = hashSetOf() + val relaysContacted = hashSetOf() + val writeRelayFreq = HashMap() + val liveRelays = hashSetOf() + + // Concurrent: touched by more than one of producer/consumer/drain-workers. + val relayHints = ConcurrentMap>() + val attempts = ConcurrentMap() + val deadRelays = ConcurrentSet() + val relayStrikes = ConcurrentMap() + + var rounds = 0 + var contactListsFed = 0 + + /** + * A relay that HARD-failed (bad domain, TLS misconfig, dead HTTP code) is + * dropped on the first strike: it will not fix itself. A TRANSIENT failure + * (refused/reset/unreachable, or a 429/5xx) might clear, so it takes + * MAX_DEAD_STRIKES before we give up. Pure timeouts never reach here — the + * drain treats them as busy-retry and does not report them dead at all. + */ + fun recordDead(failed: Map) { + for ((r, kind) in failed) { + when (kind) { + DrainFailure.HARD -> deadRelays.add(r) + DrainFailure.TRANSIENT -> + if (relayStrikes.merge(r, 1) { a, b -> a + b } >= MAX_DEAD_STRIKES) deadRelays.add(r) + } + } + } + + /** The busiest live relays we've learned, excluding the dead ones. */ + fun topLiveRelays(cap: Int): List = + writeRelayFreq.entries + .asSequence() + .filter { it.key in liveRelays && it.key !in deadRelays } + .sortedByDescending { it.value } + .take(cap) + .map { it.key } + .toList() + + /** + * Feed a user's contact list into the graph, harvest relay hints, stamp + * the hop distance of newly-seen follows, and add them to the frontier. + * Called once per user (guarded by `done`). Returns the count of + * newly-discovered users. + */ + fun ingest( + source: HexKey, + contacts: ContactListEvent, + ): Int { + val nextHop = (hopOf[source] ?: 0) + 1 + val follows = ArrayList() + var fresh = 0 + for (tag in contacts.follows()) { + follows.add(tag.pubKey) + tag.relayUri?.let { relayHints.getOrPut(tag.pubKey) { ConcurrentSet() }.add(it) } + if (discovered.add(tag.pubKey)) { + hopOf[tag.pubKey] = nextHop + fresh++ + } + } + builder.addFollows(source, follows) + contactListsFed++ + return fresh + } + + /** + * Feed into the graph the contact lists a drain just returned (deduped by + * author; the store's canonical latest wins), marking fed authors done. + * Only the authors we actually received are touched — no scan over the + * whole still-missing set. Returns the count newly fed. + */ + suspend fun harvest(events: List>): Int { + var got = 0 + for ((_, ev) in events) { + if (ev !is ContactListEvent) continue + val pk = ev.pubKey + if (pk in done) continue + val contacts = contactsOf(pk) ?: continue + done += pk + ingest(pk, contacts) + got++ + } + return got + } + + /** + * Sharded backbone sweep (see SHARD_RELAYS). Splits the missing authors + * across the top live relays — one shard per relay, so no relay gets the + * same list twice — drains all shards concurrently, then rotates whoever's + * still missing onto a different relay for up to SHARD_ROTATIONS passes. + * Once the remainder is small it's cheap to broadcast it to every top relay + * at once. Returns lists fed. + */ + suspend fun shardedSweep(authors: Collection): Int { + val top = topLiveRelays(SHARD_RELAYS) + if (top.isEmpty()) return 0 + val n = top.size + var missing = authors.filter { it !in done && contactsOf(it) == null } + var got = 0 + var rotation = 0 + while (missing.size > SHARD_BROADCAST_THRESHOLD && rotation < SHARD_ROTATIONS) { + val shards = Array(n) { ArrayList() } + for (pk in missing) { + val base = ((pk.hashCode() % n) + n) % n + shards[(base + rotation) % n].add(pk) + } + val results = + coroutineScope { + top + .mapIndexedNotNull { i, relay -> + val shard = shards[i] + if (shard.isEmpty()) { + null + } else { + // Each drain gets its own dead-set — the concurrent + // drains must not share a mutable HashMap. + async { + val dead = HashMap() + val filters = + mapOf(relay to shard.chunked(AUTHORS_PER_FILTER).map { Filter(kinds = FETCH_KINDS, authors = it) }) + drainGated(filters, dead) to dead + } + } + }.awaitAll() + } + for ((_, dead) in results) recordDead(dead) + relaysContacted += top + val flat = results.flatMap { it.first } + for ((relay, _) in flat) liveRelays.add(relay) + got += harvest(flat) + missing = missing.filter { it !in done } + rotation++ + } + // Once the remainder is small it's cheap to ask every top relay for it + // at once. If the rotations bailed with a still-large set, those authors + // just aren't on the popular relays — leave them to the caller's outbox + // pass rather than broadcast a huge list. + if (missing.isNotEmpty() && missing.size <= SHARD_BROADCAST_THRESHOLD) { + // Broadcast the small remainder to a wider set of busy relays than + // the rotation used — recovers users whose list is only on a relay + // ranked below the top SHARD_RELAYS. + val live = topLiveRelays(BROADCAST_RELAYS) + if (live.isNotEmpty()) { + val dead = HashMap() + val filters = + live.associateWith { missing.chunked(AUTHORS_PER_FILTER).map { Filter(kinds = FETCH_KINDS, authors = it) } } + val events = drainGated(filters, dead) + recordDead(dead) + relaysContacted += live + for ((relay, _) in events) liveRelays.add(relay) + got += harvest(events) + } + } + return got + } + + /** + * Fetch kind:10002 relay lists for any [pubkeys] we don't already know, so + * [routeByOutbox] can route their content query to their own write relays. + * + * Tier 1 queries the bounded relay-list discovery set (indexers + general + * defaults), which aggregate kind:10002 for the whole network. Blocking, + * because this round's routing needs the result. + * + * Tier 2 is a completeness net for the stragglers the indexers don't cover: + * cast the widest net — every relay we've seen deliver events. Fired + * fire-and-forget on [bgScope]: a stray 10002 might sit on any one relay, so + * we don't skip any, but we can't block the crawl on a fan-out that large. + * The results land in the store and improve routing for later rounds. + */ + suspend fun ensureRelayLists( + pubkeys: Set, + allLiveRelays: Set, + bgScope: CoroutineScope, + ) { + val missing = pubkeys.filter { relaysOf(it) == null } + if (missing.isEmpty()) return + + suspend fun query( + authors: List, + relays: Set, + ) { + if (relays.isEmpty() || authors.isEmpty()) return + val filters = + relays.associateWith { + authors.chunked(AUTHORS_PER_FILTER).map { chunk -> + Filter(kinds = listOf(AdvertisedRelayListEvent.KIND), authors = chunk) + } + } + drainGated(filters, null) + } + + val discovery = config.relayListDiscoveryRelays + query(missing, discovery) + + val stillMissing = missing.filter { relaysOf(it) == null } + val wide = allLiveRelays - discovery + if (stillMissing.isNotEmpty() && wide.isNotEmpty()) { + bgScope.launch { query(stillMissing, wide) } + } + } + + /** + * Fetch NIP-09 kind:5 deletion requests that retract any report we gathered. + * A reporter can delete their own kind:1984 report — a deletion valid only + * from the reporter's own key, published to the reporter's outbox. So we + * group report ids by their author and ask each author's write relays for + * kind:5 events that cite those ids (`#e`), pulling only the deletions that + * touch our reports. The events land in the store for the caller to apply. + */ + suspend fun fetchReportDeletions(backbone: Set) { + val idsByAuthor = HashMap>() + for (ev in store.query(Filter(kinds = listOf(ReportEvent.KIND)))) { + if (ev is ReportEvent) idsByAuthor.getOrPut(ev.pubKey) { ArrayList() }.add(ev.id) + } + if (idsByAuthor.isEmpty()) return + + // Route each reporter to their own write relays (fallback: backbone). + val perRelayAuthors = HashMap>() + for (author in idsByAuthor.keys) { + val write = relaysOf(author)?.writeRelaysNorm()?.takeIf { it.isNotEmpty() } ?: backbone + for (relay in write) if (relay !in deadRelays) perRelayAuthors.getOrPut(relay) { HashSet() }.add(author) + } + if (perRelayAuthors.isEmpty()) return + + val filters = + perRelayAuthors.mapValues { (_, authors) -> + buildList { + for (authorChunk in authors.chunked(AUTHORS_PER_FILTER)) { + // Scope #e to this author-chunk's own report ids, chunked to + // respect REQ limits. Any over-match (a filter pairing an + // author with another author's id) is harmless — the + // deleter-must-be-author check the caller runs rejects it. + val chunkIds = authorChunk.flatMap { idsByAuthor[it].orEmpty() } + for (idChunk in chunkIds.chunked(AUTHORS_PER_FILTER)) { + add(Filter(kinds = listOf(DeletionEvent.KIND), authors = authorChunk, tags = mapOf("e" to idChunk))) + } + } + } + } + drainGated(filters, null) + } + + /** + * Group [pubkeys] by the relays we should query for their events: + * - first try: the user's own kind:10002 write relays (the outbox model); + * - a retry (`attempts[pk] > 0`, its outbox already failed): outbox + + * [backbone] — the known-good relays other people write to; + * - no outbox at all: harvested hints + backbone + the general fallback. + * + * Also tallies each user's write relays into [writeRelayFreq] so the + * backbone can be learned from the crawl. Authors are chunked per relay. + */ + suspend fun routeByOutbox( + pubkeys: Set, + backbone: Set, + ): Map> { + val fallback = config.contentFallbackRelays + val perRelay = HashMap>() + + for (pk in pubkeys) { + val write = relaysOf(pk)?.writeRelaysNorm()?.takeIf { it.isNotEmpty() } + write?.forEach { writeRelayFreq[it] = (writeRelayFreq[it] ?: 0) + 1 } + val relays = + when { + write == null -> relayHints[pk]?.snapshot().orEmpty() + backbone + fallback + (attempts[pk] ?: 0) > 0 -> write + backbone + else -> write + } + // Skip relays already proven dead — routing to them only burns the + // drain timeout. + for (relay in relays) if (relay !in deadRelays) perRelay.getOrPut(relay) { HashSet() }.add(pk) + } + + return perRelay.mapValues { (_, authors) -> + authors.chunked(AUTHORS_PER_FILTER).map { chunk -> + Filter(kinds = FETCH_KINDS, authors = chunk) + } + } + } + + suspend fun run(): Stats { + val crawlMark = TimeSource.Monotonic.markNow() + // Scope for fire-and-forget relay-list discovery (see ensureRelayLists + // Tier 2). SupervisorJob so one failing sweep never cancels the others; + // cancelled when the crawl finishes. + val bgScope = CoroutineScope(coroutineContext + SupervisorJob()) + hopOf[observer] = 0 + + while (rounds < config.maxRounds) { + // Only crawl users within the hop budget; deeper users still appear + // in the graph as follow targets, we just don't fetch their lists. + val pending = discovered.filter { it !in done && (hopOf[it] ?: 0) < config.maxHops } + if (pending.isEmpty()) break + rounds++ + + // Refresh the warm pool to this round's busiest relays and keep that + // subscription open — reusing the same subId just updates the + // desired-relay set, so these sockets stay up across the round. + topLiveRelays(WARM_POOL_SIZE).takeIf { it.isNotEmpty() }?.let { warm -> + client.subscribe(WARM_SUB_ID, warm.associateWith { WARM_FILTERS }, null) + } + + val discoveredBefore = discovered.size + val fedBefore = contactListsFed + + // Phase A — bulk-fetch from the busiest relays via the sharded sweep. + // Most users' kind:3 lives on the big popular relays, so this clears + // the majority cheaply (early rounds no-op until a backbone is learned). + shardedSweep(pending) + + // Phase B — whoever the popular relays didn't have (niche outboxes): + // resolve their kind:10002, then fetch from their own write relays, + // drained a few at a time and skipping dead relays. + val stragglers = pending.filter { it !in done } + if (stragglers.isNotEmpty()) { + val backbone = topLiveRelays(BACKBONE_SIZE).toSet() + // Snapshot of every relay we've seen work, for the wide Tier-2 + // sweep (taken now, before the Phase-B workers mutate liveRelays). + val allLive = liveRelays.filterTo(HashSet()) { it !in deadRelays } + ensureRelayLists(stragglers.toSet(), allLive, bgScope) + + // Continuous worker pool instead of chunked awaitAll barriers, so + // no worker waits on a slow sibling and hot relays stay connected. + // Shared graph state stays single-writer: routeByOutbox runs only + // on the producer (keeps writeRelayFreq serial) and ingest runs + // only on the consumer (keeps discovered/done/builder/hopOf serial), + // now overlapped with draining instead of blocked behind each batch. + val routed = Channel, Map>>>(DRAIN_CONCURRENCY * 2) + val drainedOut = Channel, Set, List>>>(Channel.UNLIMITED) + coroutineScope { + // Producer: route each batch by outbox (serial), backpressured + // by the bounded `routed` channel. + val producer = + launch { + for (batch in stragglers.chunked(USER_BATCH)) { + val filters = routeByOutbox(batch.toSet(), backbone) + routed.send(batch to filters) + } + routed.close() + } + // Drain workers: pure network, no shared graph-state writes + // except recordDead (concurrent-safe). + val workers = + List(DRAIN_CONCURRENCY) { + launch { + for ((batch, filters) in routed) { + val dead = HashMap() + val events = drainGated(filters, dead) + recordDead(dead) + drainedOut.send(Triple(batch, filters.keys, events)) + } + } + } + // Consumer: single-writer ingest, overlapped with draining. + val consumer = + launch { + for ((batch, relays, events) in drainedOut) { + relaysContacted += relays + // Any relay that gave us an event is proven live + useful. + for ((relay, _) in events) liveRelays.add(relay) + for (pk in batch) { + if (pk in done) continue + val contacts = contactsOf(pk) + if (contacts != null) { + done += pk + ingest(pk, contacts) + } else { + val tries = (attempts[pk] ?: 0) + 1 + attempts[pk] = tries + if (tries >= MAX_OUTBOX_ATTEMPTS) done += pk + } + } + } + } + producer.join() + workers.joinAll() + drainedOut.close() + consumer.join() + } + } + + log( + "[graperank] round $rounds: pending=${pending.size}, " + + "gotList=${contactListsFed - fedBefore}, newUsers=${discovered.size - discoveredBefore}, " + + "discovered=${discovered.size}, done=${done.size}, dead=${deadRelays.size()}", + ) + } + + // Crawl done — drop the warm pool and stop any background relay-list + // sweeps still in flight (their results are already in the store). + client.unsubscribe(WARM_SUB_ID) + bgScope.cancel() + + // Reports can be retracted. Ask each reporter's outbox for NIP-09 kind:5 + // deletions that cite the reports we gathered (#e-filtered to our report + // ids). The events land in the store; the caller decides which reports + // they actually retract. + fetchReportDeletions(topLiveRelays(BACKBONE_SIZE).toSet()) + + val hopHistogram = + hopOf.values + .groupingBy { it } + .eachCount() + .toList() + .sortedBy { it.first } + .toMap() + val downloadMs = crawlMark.elapsedNow().inWholeMilliseconds + log( + "[graperank] crawl complete: ${discovered.size} discovered, $contactListsFed contact lists fed, " + + "${relaysContacted.size} relays contacted, ${deadRelays.size()} dead, $rounds rounds in $downloadMs ms; " + + "by hop: " + hopHistogram.entries.joinToString(" ") { "${it.key}=${it.value}" }, + ) + return Stats( + rounds = rounds, + discovered = discovered.size, + contactListsFed = contactListsFed, + relaysContacted = relaysContacted.size, + deadRelays = deadRelays.size(), + hopHistogram = hopHistogram, + downloadMs = downloadMs, + ) + } + } + + /** + * Subscribe each relay to its filters behind [limiter], drain until every + * relay's subscription is terminal or the timeout elapses, verify+store the + * events, and return them tagged by relay. Each relay gets its own gated + * subscription so we never exceed its adaptive concurrent-subscription cap; a + * relay's filters are split into REQ-sized groups so a popular relay routed + * thousands of authors doesn't produce a multi-MB frame that most relays + * reject outright. Hard connect failures are reported into [deadOut]. + */ + private suspend fun drainGated( + filters: Map>, + deadOut: MutableMap?, + ): List> { + if (filters.isEmpty()) return emptyList() + val eventChannel = Channel>(Channel.UNLIMITED) + + // Split each relay's filters into REQ-sized groups. A REQ frame carries ALL + // its filters at once, so a popular relay routed thousands of authors would + // otherwise produce a multi-MB frame that most relays reject ("message too + // large") — silently dropping every author in it. Grouping by total entry + // count keeps each REQ well under the common 256KB cap. + val units = ArrayList>>() + for ((relay, relayFilters) in filters) { + var group = ArrayList() + var entries = 0 + for (f in relayFilters) { + val fe = filterEntries(f) + if (group.isNotEmpty() && entries + fe > MAX_REQ_ENTRIES) { + units.add(relay to group) + group = ArrayList() + entries = 0 + } + group.add(f) + entries += fe + } + if (group.isNotEmpty()) units.add(relay to group) + } + + // Per-relay failure classification, HARD winning over TRANSIENT across a + // relay's several REQ-groups; plus which relays stalled to a timeout. + val failures = ConcurrentMap() + val timedOut = ConcurrentSet() + + val collected = mutableListOf>() + coroutineScope { + // Single consumer: verify+store serially. One writer, so SeenIds' + // single-writer contract holds. The outbox model delivers the SAME event + // from many relays at once; skip a duplicate BEFORE the expensive Schnorr + // verify+store. An id is marked seen only after it verifies, so a forged + // copy (valid id, bad signature) delivered first can't suppress the + // genuine one that follows. + val consumer = + launch { + val seen = SeenIds(initialSlotsPow2 = 12) + for ((relay, event) in eventChannel) { + if (seen.contains(event.id)) continue + if (verifyAndStore(event)) { + seen.add(event.id) + collected.add(relay to event) + } + } + } + // One gated subscription per (relay, REQ-group). The permit is held for + // the group's whole life, so concurrent subs on a relay never exceed its + // adaptive cap. + units + .map { (subRelay, groupFilters) -> + launch { + limiter.withPermit(subRelay) { + val subId = newSubId() + val done = CompletableDeferred() + val groupListener = + object : SubscriptionListener { + override fun onEvent( + event: Event, + isLive: Boolean, + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + eventChannel.trySend(relay to event) + } + + override fun onEose( + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + done.complete("eose") + } + + override fun onClosed( + message: String, + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + done.complete("closed:$message") + } + + override fun onCannotConnect( + relay: NormalizedRelayUrl, + message: String, + forFilters: List?, + ) { + done.complete("cannot:$message") + } + } + client.subscribe(subId, mapOf(subRelay to groupFilters), groupListener) + try { + val reason = withTimeoutOrNull(config.timeoutMs) { done.await() } ?: "timeout" + if (reason == "timeout") timedOut.add(subRelay) + classifyDrainFailure(reason)?.let { kind -> + failures.merge(subRelay, kind) { a, b -> + if (a == DrainFailure.HARD || b == DrainFailure.HARD) DrainFailure.HARD else DrainFailure.TRANSIENT + } + } + } finally { + client.unsubscribe(subId) + } + } + } + }.joinAll() + // All subscriptions are torn down; no more events can arrive. Close the + // channel so the consumer drains what's buffered and completes. + eventChannel.close() + consumer.join() + } + if (config.diagnose && timedOut.size() > 0) { + val stalled = timedOut.snapshot() + val eventsPer = collected.groupingBy { it.first }.eachCount() + val detail = stalled.take(12).joinToString(", ") { "${it.url}(${eventsPer[it] ?: 0}ev)" } + log("[drain] timeout ${config.timeoutMs}ms: ${stalled.size} slow(no EOSE)" + (if (detail.isNotEmpty()) " | slow: $detail" else "")) + } + deadOut?.putAll(failures.snapshot()) + return collected + } + + /** + * Verify [event]'s NIP-01 id+signature and, if valid, persist it to [store]. + * Returns true when the event was accepted. A UNIQUE-constraint rejection is + * normal (the store already holds this id, or a newer replaceable) — the outbox + * model delivers the same event from several relays, so a crawl produces these + * by the hundred-thousand — so only genuine persistence failures are logged. + */ + private suspend fun verifyAndStore(event: Event): Boolean { + if (!event.verify()) { + Log.w("GrapeRankDataCrawler") { "dropped event ${event.id.take(8)} kind=${event.kind} — bad signature" } + return false + } + try { + store.insert(event) + } catch (t: Throwable) { + if (t.message?.contains("UNIQUE constraint", ignoreCase = true) != true) { + Log.w("GrapeRankDataCrawler") { "store insert failed for ${event.id.take(8)}: ${t.message}" } + } + } + return true + } + + /** Latest known kind:3 contact list for [pubKey] from the local store, or null. */ + private suspend fun contactsOf(pubKey: HexKey): ContactListEvent? = + store + .query(Filter(authors = listOf(pubKey), kinds = listOf(ContactListEvent.KIND), limit = 1)) + .firstOrNull() as? ContactListEvent + + /** Latest known kind:10002 advertised relay list for [pubKey] from the store, or null. */ + private suspend fun relaysOf(pubKey: HexKey): AdvertisedRelayListEvent? = + store + .query(Filter(authors = listOf(pubKey), kinds = listOf(AdvertisedRelayListEvent.KIND), limit = 1)) + .firstOrNull() as? AdvertisedRelayListEvent + + companion object { + // Authors per REQ filter — keeps individual subscriptions within relay limits. + private const val AUTHORS_PER_FILTER = 300 + + // Max total "entries" (authors + ids + tag values) in a single REQ frame. + // Each entry is a ~67-byte hex string, so 2500 ≈ 167KB — under the 256KB + // message cap most relays enforce. drainGated groups filters to stay within. + private const val MAX_REQ_ENTRIES = 2500 + + // Times we re-query an unreachable user's outbox before giving up, so the + // crawl still terminates on a finite graph. + private const val MAX_OUTBOX_ATTEMPTS = 3 + + // Users whose outboxes we fetch in a single drain. Draining thousands of + // distinct outbox relays at once saturates connections and times out + // (~250/drain succeeds, ~17k fails); keep the fan-out small. + private const val USER_BATCH = 256 + + // Global content-drain fan-out — how many outbox batches we drain at once. A + // GLOBAL bound (memory / open sockets); the per-relay concurrency limit is + // enforced separately by AdaptiveRelayLimiter. A higher global fan-out + // re-floods busy hubs faster than demotion catches up, so keep it moderate. + private const val DRAIN_CONCURRENCY = 24 + + // Sharded backbone sweep: split the still-missing authors into SHARD_RELAYS + // lists, one per top relay, rotating up to SHARD_ROTATIONS times; once the + // remainder drops below SHARD_BROADCAST_THRESHOLD, broadcast it at once. + private const val SHARD_RELAYS = 10 + private const val SHARD_ROTATIONS = 6 + private const val SHARD_BROADCAST_THRESHOLD = 2000 + + // The small-remainder broadcast goes to this many top live relays — a user's + // kind:3 is often mirrored on a busy relay ranked below the top 10. + private const val BROADCAST_RELAYS = 60 + + // A relay that fails to CONNECT this many times is treated as dead. Kept + // above 1 so a single transient connect blip doesn't evict a relay. + private const val MAX_DEAD_STRIKES = 3 + + // Most-used write relays kept as the known-good backbone for retrying users. + private const val BACKBONE_SIZE = 30 + + // Warm pool: hold a do-nothing subscription open to the busiest relays for + // the whole crawl, so the connections we reuse every round survive the + // between-round routing gaps. The filter matches an impossible event id, so + // the relay EOSEs immediately and streams nothing — it only keeps sockets warm. + private const val WARM_POOL_SIZE = 20 + private const val WARM_SUB_ID = "graperank-warm" + private val WARM_FILTERS = listOf(Filter(ids = listOf("0".repeat(64)))) + + // Kinds requested from relays during the crawl: the graph edges (contact + // lists, mute lists, reports) PLUS the user's own kind:10002. A user's outbox + // holds the freshest copy of their relay list, so folding 10002 into the same + // query keeps routing current. The store keeps newest-by-created_at for the + // replaceable 10002, so the freshest always wins regardless of source relay. + private val FETCH_KINDS = + listOf(ContactListEvent.KIND, MuteListEvent.KIND, ReportEvent.KIND, AdvertisedRelayListEvent.KIND) + + /** Count the size-driving entries in a filter: authors, ids, and tag values. */ + private fun filterEntries(f: Filter): Int = + (f.authors?.size ?: 0) + + (f.ids?.size ?: 0) + + (f.tags?.values?.sumOf { it.size } ?: 0) + + (f.tagsAll?.values?.sumOf { it.size } ?: 0) + } +} diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/AdaptiveRelayLimiter.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/AdaptiveRelayLimiter.kt similarity index 80% rename from cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/AdaptiveRelayLimiter.kt rename to quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/AdaptiveRelayLimiter.kt index b2f39fc534..c66134480f 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/AdaptiveRelayLimiter.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/AdaptiveRelayLimiter.kt @@ -18,7 +18,7 @@ * 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.amethyst.cli +package com.vitorpamplona.quartz.nip01Core.relay.client.accessories import com.vitorpamplona.quartz.nip01Core.relay.client.listeners.RelayConnectionListener import com.vitorpamplona.quartz.nip01Core.relay.client.single.IRelayClient @@ -26,13 +26,16 @@ import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.ClosedMessage import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.NoticeMessage import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import com.vitorpamplona.quartz.utils.Log +import com.vitorpamplona.quartz.utils.TimeUtils +import com.vitorpamplona.quartz.utils.concurrent.ConcurrentMap import kotlinx.coroutines.CompletableDeferred import kotlinx.coroutines.delay import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock -import java.util.concurrent.ConcurrentHashMap -import java.util.concurrent.atomic.AtomicInteger -import java.util.concurrent.atomic.AtomicLong +import kotlin.concurrent.atomics.AtomicInt +import kotlin.concurrent.atomics.AtomicLong +import kotlin.concurrent.atomics.ExperimentalAtomicApi /** * Adaptive per-relay back-pressure with TWO independent controls, because relays @@ -60,26 +63,27 @@ import java.util.concurrent.atomic.AtomicLong * Registered as a [RelayConnectionListener] on the shared client, so both signals * are driven straight off the incoming NOTICE/CLOSED frames (which fire on the * per-relay socket threads — all state here is concurrent). Drains gate through - * [withPermit]; [Context.drain]'s `gatePerRelay` path holds a relay's permit for - * the lifetime of that relay's subscription, and passes the rate gate before it - * opens, so we respect both limits at once. + * [withPermit]: the gated-drain path holds a relay's permit for the lifetime of + * that relay's subscription, and passes the rate gate before it opens, so we + * respect both limits at once. */ +@OptIn(ExperimentalAtomicApi::class) class AdaptiveRelayLimiter( private val startCap: Int = 100, private val subLadder: List = listOf(20, 10), private val rateLadder: List = listOf(250L, 500L, 1000L, 2000L), ) : RelayConnectionListener { - private val gates = ConcurrentHashMap() + private val gates = ConcurrentMap() // Concurrency-cap demotions per relay (== index+1 into subLadder). Capped at // subLadder.size: past the floor we stop demoting. - private val subDemotions = ConcurrentHashMap() + private val subDemotions = ConcurrentMap() // Rate-limit state per relay: how far down rateLadder we've stepped, the // current min interval between opens, and the next epoch-ms an open may fire. - private val rateSteps = ConcurrentHashMap() - private val rateDelayMs = ConcurrentHashMap() - private val nextAllowedAtMs = ConcurrentHashMap() + private val rateSteps = ConcurrentMap() + private val rateDelayMs = ConcurrentMap() + private val nextAllowedAtMs = ConcurrentMap() private fun gate(relay: NormalizedRelayUrl): Gate = gates.getOrPut(relay) { Gate(startCap) } @@ -106,14 +110,14 @@ class AdaptiveRelayLimiter( private suspend fun rateGate(relay: NormalizedRelayUrl) { val delayMs = rateDelayMs[relay] ?: return if (delayMs <= 0L) return - val now = System.currentTimeMillis() + val now = TimeUtils.nowMillis() // Atomically claim the next slot: my turn is max(prevSlot, now); the next // caller can't fire until delayMs after me. Serializes opens to this relay // at one per delayMs, in arrival order. val slot = nextAllowedAtMs.getOrPut(relay) { AtomicLong(now) } var myTurn: Long while (true) { - val prev = slot.get() + val prev = slot.load() myTurn = maxOf(prev, now) if (slot.compareAndSet(prev, myTurn + delayMs)) break } @@ -142,49 +146,51 @@ class AdaptiveRelayLimiter( /** Step [relay] one rung down the concurrency-cap ladder, unless already at the floor. */ private fun demoteConcurrency(relay: NormalizedRelayUrl) { if ((subDemotions[relay] ?: 0) >= subLadder.size) return - val step = subDemotions.merge(relay, 1, Int::plus)!! + val step = subDemotions.merge(relay, 1) { a, b -> a + b } val cap = subLadder[(step - 1).coerceIn(0, subLadder.size - 1)] gate(relay).lower(cap) if (step <= subLadder.size) { - System.err.println("[limiter] ${relay.url} concurrency capped at $cap subs (sub-limit #$step)") + Log.w("AdaptiveRelayLimiter") { "${relay.url} concurrency capped at $cap subs (sub-limit #$step)" } } } /** Step [relay] one rung down the rate ladder, unless already at the slowest. */ private fun throttleRate(relay: NormalizedRelayUrl) { if ((rateSteps[relay] ?: 0) >= rateLadder.size) return - val step = rateSteps.merge(relay, 1, Int::plus)!! + val step = rateSteps.merge(relay, 1) { a, b -> a + b } val d = rateLadder[(step - 1).coerceIn(0, rateLadder.size - 1)] rateDelayMs[relay] = d if (step <= rateLadder.size) { - System.err.println("[limiter] ${relay.url} rate-throttled to 1 REQ / ${d}ms (rate-limit #$step)") + Log.w("AdaptiveRelayLimiter") { "${relay.url} rate-throttled to 1 REQ / ${d}ms (rate-limit #$step)" } } } /** JSON-friendly view of which relays we throttled, in which dimension, how far. */ fun snapshot(): Map { - val cappedAt = sortedMapOf() - for ((_, step) in subDemotions) { + val capCounts = HashMap() + for ((_, step) in subDemotions.snapshot()) { val cap = subLadder[(step - 1).coerceIn(0, subLadder.size - 1)] - cappedAt.merge(cap, 1, Int::plus) + capCounts[cap] = (capCounts[cap] ?: 0) + 1 } - val rateAt = sortedMapOf() - for ((_, step) in rateSteps) { + val cappedAt = capCounts.toList().sortedBy { it.first }.toMap() + val rateCounts = HashMap() + for ((_, step) in rateSteps.snapshot()) { val d = rateLadder[(step - 1).coerceIn(0, rateLadder.size - 1)] - rateAt.merge(d, 1, Int::plus) + rateCounts[d] = (rateCounts[d] ?: 0) + 1 } + val rateAt = rateCounts.toList().sortedBy { it.first }.toMap() return mapOf( "start_cap" to startCap, "sub_ladder" to subLadder, "rate_ladder_ms" to rateLadder, - "concurrency_capped_relays" to subDemotions.size, + "concurrency_capped_relays" to subDemotions.size(), "concurrency_capped_at" to cappedAt, - "rate_limited_relays" to rateSteps.size, + "rate_limited_relays" to rateSteps.size(), "rate_limited_at_ms" to rateAt, ) } - fun hadThrottling(): Boolean = subDemotions.isNotEmpty() || rateSteps.isNotEmpty() + fun hadThrottling(): Boolean = subDemotions.size() > 0 || rateSteps.size() > 0 /** * A bounded-concurrency gate whose limit can only ever be *lowered* (relays @@ -197,7 +203,7 @@ class AdaptiveRelayLimiter( private class Gate( initialLimit: Int, ) { - private val limit = AtomicInteger(initialLimit) + private val limit = AtomicInt(initialLimit) private val mutex = Mutex() private var inUse = 0 private val waiters = ArrayDeque>() @@ -205,7 +211,7 @@ class AdaptiveRelayLimiter( suspend fun acquire() { val wait = mutex.withLock { - if (inUse < limit.get()) { + if (inUse < limit.load()) { inUse++ null } else { @@ -218,7 +224,7 @@ class AdaptiveRelayLimiter( suspend fun release() { mutex.withLock { inUse-- - while (inUse < limit.get() && waiters.isNotEmpty()) { + while (inUse < limit.load() && waiters.isNotEmpty()) { waiters.removeFirst().complete(Unit) inUse++ } @@ -227,7 +233,11 @@ class AdaptiveRelayLimiter( /** Monotonically shrink the cap. Safe to call from any thread. */ fun lower(newLimit: Int) { - limit.updateAndGet { if (newLimit < it) newLimit else it } + while (true) { + val cur = limit.load() + if (newLimit >= cur) return + if (limit.compareAndSet(cur, newLimit)) return + } } } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/DrainFailure.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/DrainFailure.kt new file mode 100644 index 0000000000..b2d605d044 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/DrainFailure.kt @@ -0,0 +1,77 @@ +/* + * 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 + +/** + * Why a relay could not be used for a one-shot drain — when the reason is worth + * acting on (dropping the relay from further routing). + * + * - [HARD]: the relay answered wrong, or cannot exist. A bad HTTP upgrade (not a + * websocket / dead status code), an unresolvable domain, or a TLS misconfig. + * This will not fix itself, so one strike is enough to drop it. + * - [TRANSIENT]: a failure that might clear — connection refused / reset, host + * unreachable, or a temporary 429/5xx on the upgrade. Struck a few times + * before we give up. + * + * A pure connect **timeout** is neither. The relay is most likely just busy, so + * we retry it and never mark it dead — [classifyDrainFailure] returns null for + * it (and for any non-failure terminal reason). + */ +enum class DrainFailure { HARD, TRANSIENT } + +/** + * Classify a drain per-relay terminal reason. Returns null when the relay should + * simply be retried (a timeout, or a non-failure like eose/closed). The reason + * shape is `cannot:` for a connect failure (see + * `BasicRelayClient.onCannotConnect`), or `eose` / `closed:…` / `timeout`. + */ +fun classifyDrainFailure(reason: String): DrainFailure? { + if (!reason.startsWith("cannot")) return null + val m = reason.removePrefix("cannot:").lowercase() + // The message now carries the exception class name (see BasicRelayClient), so + // we can key on the stable *type* rather than localized message text. + // Busy, not dead: a connect/read timeout means the handshake just didn't + // finish in time. Retry it — the relay is probably fine, only slow or loaded. + if ("timeout" in m || "timed out" in m) return null // SocketTimeoutException, etc. + // Cannot ever work: unresolvable domain (DNS) or a TLS misconfiguration. + // Dead for good — one strike is enough. + if ("unknownhost" in m || // UnknownHostException + "unable to resolve host" in m || + "no address associated" in m || + "nodename nor servname" in m || + "sslhandshake" in m || // SSLHandshakeException + "sslpeerunverified" in m || + "sslexception" in m || + "certificate" in m || // CertificateException + "trust anchor" in m || + "certpath" in m + ) { + return DrainFailure.HARD + } + // Wrong HTTP upgrade. Usually a misconfigured endpoint (not a relay), but + // 429 / 5xx mean "busy, come back later", so those stay transient. + if ("server misconfigured" in m || "not a websocket" in m || "expected http 101" in m) { + val transientCode = Regex("response: (429|500|502|503|504)").containsMatchIn(m) + return if (transientCode) DrainFailure.TRANSIENT else DrainFailure.HARD + } + // Refused / reset / unreachable / anything else: might clear — retry a few times. + return DrainFailure.TRANSIENT +} diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.kt new file mode 100644 index 0000000000..47d3233c9b --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.kt @@ -0,0 +1,68 @@ +/* + * 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.utils.concurrent + +/** + * A thread-safe hash map whose compound operations — [getOrPut] and [merge] — + * apply their update **atomically**, not merely one-lock-per-primitive-op. This + * is the contract a concurrent producer/consumer pipeline needs: two coroutines + * racing `getOrPut` on the same key must agree on a single value, and racing + * `merge` must not lose an increment. + * + * commonMain has no `java.util.concurrent.ConcurrentHashMap`, so this is + * expect/actual, matching the split already used by [com.vitorpamplona.quartz.utils.cache.ConcurrentHashCache]: + * - JVM / Android → `ConcurrentHashMap` (lock-free, true atomic `computeIfAbsent` / `merge`). + * - Native (Apple + Linux) → copy-on-write over an atomic reference, with a + * CAS retry loop giving the same atomicity. Correct but O(n)-per-write; the + * native targets never run the heavy crawl this backs, they only compile it. + * + * Only the operations the crawl actually uses are exposed — no full [MutableMap] + * surface — so the native copy-on-write actual stays small and obviously correct. + */ +expect class ConcurrentMap() { + operator fun get(key: K): V? + + operator fun set( + key: K, + value: V, + ) + + /** Atomically return the value for [key], computing and inserting [defaultValue] once if absent. */ + fun getOrPut( + key: K, + defaultValue: () -> V, + ): V + + /** + * Atomically insert [value] if [key] is absent, else replace the existing + * value with `remap(existing, value)`. Returns the value now stored. + */ + fun merge( + key: K, + value: V, + remap: (old: V, new: V) -> V, + ): V + + fun size(): Int + + /** A point-in-time copy of the entries — safe to iterate without holding a lock. */ + fun snapshot(): Map +} diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentSet.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentSet.kt new file mode 100644 index 0000000000..d00545acf6 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentSet.kt @@ -0,0 +1,43 @@ +/* + * 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.utils.concurrent + +/** + * A thread-safe hash set for the crawl's cross-coroutine membership tracking + * (dead relays struck by drain workers while the router reads them, relay hints + * written by the ingest consumer while the producer reads them). + * + * commonMain has no `java.util.concurrent.ConcurrentHashMap.newKeySet()`, so this + * is expect/actual with the same JVM-vs-native split as [ConcurrentMap]: + * - JVM / Android → `ConcurrentHashMap.newKeySet()`. + * - Native → copy-on-write over an atomic reference (compile-only, never the hot path). + */ +expect class ConcurrentSet() { + /** Add [element]; returns true if it was not already present. */ + fun add(element: E): Boolean + + operator fun contains(element: E): Boolean + + fun size(): Int + + /** A point-in-time copy — safe to iterate or diff against without a lock. */ + fun snapshot(): Set +} diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentCollectionsTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentCollectionsTest.kt new file mode 100644 index 0000000000..208054f818 --- /dev/null +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentCollectionsTest.kt @@ -0,0 +1,108 @@ +/* + * 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.utils.concurrent + +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFalse +import kotlin.test.assertNull +import kotlin.test.assertTrue + +class ConcurrentCollectionsTest { + @Test + fun mapGetSet() { + val m = ConcurrentMap() + assertNull(m["a"]) + m["a"] = 1 + assertEquals(1, m["a"]) + m["a"] = 2 + assertEquals(2, m["a"]) + assertEquals(1, m.size()) + } + + @Test + fun mapGetOrPutComputesOnce() { + val m = ConcurrentMap() + var calls = 0 + assertEquals( + 7, + m.getOrPut("k") { + calls++ + 7 + }, + ) + // Present now: the default must NOT be recomputed. + assertEquals( + 7, + m.getOrPut("k") { + calls++ + 99 + }, + ) + assertEquals(1, calls) + assertEquals(7, m["k"]) + } + + @Test + fun mapMergeInsertsThenCombines() { + val m = ConcurrentMap() + // Absent -> inserts the value verbatim, remap not applied. + assertEquals(1, m.merge("k", 1) { a, b -> a + b }) + // Present -> remap(existing, value). + assertEquals(4, m.merge("k", 3) { a, b -> a + b }) + assertEquals(4, m["k"]) + } + + @Test + fun mapSnapshotIsDetached() { + val m = ConcurrentMap() + m["a"] = 1 + m["b"] = 2 + val snap = m.snapshot() + assertEquals(mapOf("a" to 1, "b" to 2), snap) + // Mutating the map after the snapshot must not change the snapshot. + m["c"] = 3 + assertEquals(2, snap.size) + assertEquals(3, m.size()) + } + + @Test + fun setAddContainsSize() { + val s = ConcurrentSet() + assertFalse("x" in s) + assertTrue(s.add("x")) + // Re-adding is a no-op and reports it. + assertFalse(s.add("x")) + assertTrue("x" in s) + assertTrue(s.add("y")) + assertEquals(2, s.size()) + } + + @Test + fun setSnapshotIsDetached() { + val s = ConcurrentSet() + s.add("a") + val snap = s.snapshot() + s.add("b") + assertEquals(setOf("a"), snap) + assertEquals(2, s.size()) + } +} diff --git a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.jvmAndroid.kt b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.jvmAndroid.kt new file mode 100644 index 0000000000..65a734f8d2 --- /dev/null +++ b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.jvmAndroid.kt @@ -0,0 +1,51 @@ +/* + * 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.utils.concurrent + +import java.util.concurrent.ConcurrentHashMap + +actual class ConcurrentMap { + private val map = ConcurrentHashMap() + + actual operator fun get(key: K): V? = map[key] + + actual operator fun set( + key: K, + value: V, + ) { + map[key] = value + } + + actual fun getOrPut( + key: K, + defaultValue: () -> V, + ): V = map.computeIfAbsent(key) { defaultValue() } + + actual fun merge( + key: K, + value: V, + remap: (old: V, new: V) -> V, + ): V = map.merge(key, value) { old, new -> remap(old, new) }!! + + actual fun size(): Int = map.size + + actual fun snapshot(): Map = HashMap(map) +} diff --git a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentSet.jvmAndroid.kt b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentSet.jvmAndroid.kt new file mode 100644 index 0000000000..94a0754c11 --- /dev/null +++ b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentSet.jvmAndroid.kt @@ -0,0 +1,35 @@ +/* + * 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.utils.concurrent + +import java.util.concurrent.ConcurrentHashMap + +actual class ConcurrentSet { + private val set: MutableSet = ConcurrentHashMap.newKeySet() + + actual fun add(element: E): Boolean = set.add(element) + + actual operator fun contains(element: E): Boolean = set.contains(element) + + actual fun size(): Int = set.size + + actual fun snapshot(): Set = HashSet(set) +} diff --git a/quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.native.kt b/quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.native.kt new file mode 100644 index 0000000000..de4ee540d8 --- /dev/null +++ b/quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.native.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.utils.concurrent + +import kotlin.concurrent.atomics.AtomicReference +import kotlin.concurrent.atomics.ExperimentalAtomicApi + +// Copy-on-write, mirroring ConcurrentHashCache.linux: correct and simple. The +// native targets never run the crawl this backs (it is JVM/Android-only work); +// they only compile it, so the O(n)-per-write cost is irrelevant. A CAS retry +// loop gives getOrPut/merge the same atomicity the JVM actual gets for free. +@OptIn(ExperimentalAtomicApi::class) +actual class ConcurrentMap { + private val ref = AtomicReference(HashMap()) + + actual operator fun get(key: K): V? = ref.load()[key] + + actual operator fun set( + key: K, + value: V, + ) { + while (true) { + val cur = ref.load() + val copy = HashMap(cur) + copy[key] = value + if (ref.compareAndSet(cur, copy)) return + } + } + + actual fun getOrPut( + key: K, + defaultValue: () -> V, + ): V { + while (true) { + val cur = ref.load() + cur[key]?.let { return it } + val value = defaultValue() + val copy = HashMap(cur) + copy[key] = value + if (ref.compareAndSet(cur, copy)) return value + } + } + + actual fun merge( + key: K, + value: V, + remap: (old: V, new: V) -> V, + ): V { + while (true) { + val cur = ref.load() + val old = cur[key] + val merged = if (old == null) value else remap(old, value) + val copy = HashMap(cur) + copy[key] = merged + if (ref.compareAndSet(cur, copy)) return merged + } + } + + actual fun size(): Int = ref.load().size + + actual fun snapshot(): Map = HashMap(ref.load()) +} diff --git a/quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentSet.native.kt b/quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentSet.native.kt new file mode 100644 index 0000000000..70bf86f133 --- /dev/null +++ b/quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentSet.native.kt @@ -0,0 +1,46 @@ +/* + * 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.utils.concurrent + +import kotlin.concurrent.atomics.AtomicReference +import kotlin.concurrent.atomics.ExperimentalAtomicApi + +// Copy-on-write native actual — see ConcurrentMap.native for the rationale. +@OptIn(ExperimentalAtomicApi::class) +actual class ConcurrentSet { + private val ref = AtomicReference(HashSet()) + + actual fun add(element: E): Boolean { + while (true) { + val cur = ref.load() + if (element in cur) return false + val copy = HashSet(cur) + copy.add(element) + if (ref.compareAndSet(cur, copy)) return true + } + } + + actual operator fun contains(element: E): Boolean = element in ref.load() + + actual fun size(): Int = ref.load().size + + actual fun snapshot(): Set = HashSet(ref.load()) +}