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.
This commit is contained in:
Claude
2026-07-07 15:41:31 +00:00
parent 6bbc1b7d22
commit 667358a06a
12 changed files with 1406 additions and 843 deletions
@@ -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:<message>` 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<NormalizedRelayUrl, DrainFailure>? = null,
gatePerRelay: Boolean = false,
): List<Pair<NormalizedRelayUrl, Event>> {
if (filters.isEmpty()) return emptyList()
if (gatePerRelay) return drainGated(filters, timeoutMs, diagnoseSlow, deadOut)
val eventChannel = Channel<Pair<NormalizedRelayUrl, Event>>(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<NormalizedRelayUrl, List<Filter>>,
timeoutMs: Long,
diagnoseSlow: Boolean,
deadOut: MutableMap<NormalizedRelayUrl, DrainFailure>?,
): List<Pair<NormalizedRelayUrl, Event>> {
val eventChannel = Channel<Pair<NormalizedRelayUrl, Event>>(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<Pair<NormalizedRelayUrl, List<Filter>>>()
for ((relay, relayFilters) in filters) {
var group = ArrayList<Filter>()
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<NormalizedRelayUrl, DrainFailure>()
val timedOut = ConcurrentHashMap.newKeySet<NormalizedRelayUrl>()
val collected = mutableListOf<Pair<NormalizedRelayUrl, Event>>()
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<String>()
val groupListener =
object : SubscriptionListener {
override fun onEvent(
event: Event,
isLive: Boolean,
r: NormalizedRelayUrl,
forFilters: List<Filter>?,
) {
eventChannel.trySend(r to event)
}
override fun onEose(
r: NormalizedRelayUrl,
forFilters: List<Filter>?,
) {
done.complete("eose")
}
override fun onClosed(
message: String,
r: NormalizedRelayUrl,
forFilters: List<Filter>?,
) {
done.complete("closed:$message")
}
override fun onCannotConnect(
r: NormalizedRelayUrl,
message: String,
forFilters: List<Filter>?,
) {
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
@@ -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<NormalizedRelayUrl> =
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<String>,
@@ -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<HexKey, Int>()
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<HexKey, MutableSet<NormalizedRelayUrl>>()
// 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<HexKey>()
// Outbox retry counts. Concurrent: the producer reads (to widen a
// retry's routing) while the consumer increments.
val attempts = ConcurrentHashMap<HexKey, Int>()
val relaysContacted = hashSetOf<NormalizedRelayUrl>()
// 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<NormalizedRelayUrl, Int>()
val liveRelays = hashSetOf<NormalizedRelayUrl>()
// 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<NormalizedRelayUrl>()
val relayStrikes = ConcurrentHashMap<NormalizedRelayUrl, Int>()
// 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<NormalizedRelayUrl, DrainFailure>) {
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<NormalizedRelayUrl> =
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<HexKey>()
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<Pair<NormalizedRelayUrl, Event>>): 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<HexKey>): 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<HexKey>() }
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<NormalizedRelayUrl, DrainFailure>()
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<NormalizedRelayUrl, DrainFailure>()
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<Pair<List<HexKey>, Map<NormalizedRelayUrl, List<Filter>>>>(DRAIN_CONCURRENCY * 2)
val drainedOut = Channel<Triple<List<HexKey>, Set<NormalizedRelayUrl>, List<Pair<NormalizedRelayUrl, Event>>>>(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<NormalizedRelayUrl, DrainFailure>()
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<String, Any?>(
"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<String, Int>()),
"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<NormalizedRelayUrl> = 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<NormalizedRelayUrl> = 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<HexKey>,
allLiveRelays: Set<NormalizedRelayUrl>,
bgScope: CoroutineScope,
timeoutMs: Long,
diagnose: Boolean,
) {
val missing = pubkeys.filter { ctx.relaysOf(it) == null }
if (missing.isEmpty()) return
suspend fun query(
authors: List<HexKey>,
relays: Set<NormalizedRelayUrl>,
) {
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<NormalizedRelayUrl>,
deadRelays: Set<NormalizedRelayUrl>,
timeoutMs: Long,
diagnose: Boolean,
) {
val idsByAuthor = HashMap<HexKey, MutableList<HexKey>>()
for (ev in ctx.store.query<Event>(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<NormalizedRelayUrl, MutableSet<HexKey>>()
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<HexKey>,
hints: Map<HexKey, Set<NormalizedRelayUrl>>,
backbone: Set<NormalizedRelayUrl>,
attempts: Map<HexKey, Int>,
writeRelayFreq: MutableMap<NormalizedRelayUrl, Int>,
kinds: List<Int>,
deadRelays: Set<NormalizedRelayUrl>,
): Map<NormalizedRelayUrl, List<Filter>> {
val fallback = contentFallbackRelays(ctx)
val perRelay = HashMap<NormalizedRelayUrl, MutableSet<HexKey>>()
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
@@ -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<NormalizedRelayUrl>,
val contentFallbackRelays: Set<NormalizedRelayUrl>,
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<Int, Int>,
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<HexKey, Int>()
val discovered = hashSetOf(observer)
val done = hashSetOf<HexKey>()
val relaysContacted = hashSetOf<NormalizedRelayUrl>()
val writeRelayFreq = HashMap<NormalizedRelayUrl, Int>()
val liveRelays = hashSetOf<NormalizedRelayUrl>()
// Concurrent: touched by more than one of producer/consumer/drain-workers.
val relayHints = ConcurrentMap<HexKey, ConcurrentSet<NormalizedRelayUrl>>()
val attempts = ConcurrentMap<HexKey, Int>()
val deadRelays = ConcurrentSet<NormalizedRelayUrl>()
val relayStrikes = ConcurrentMap<NormalizedRelayUrl, Int>()
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<NormalizedRelayUrl, DrainFailure>) {
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<NormalizedRelayUrl> =
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<HexKey>()
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<Pair<NormalizedRelayUrl, Event>>): 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<HexKey>): 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<HexKey>() }
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<NormalizedRelayUrl, DrainFailure>()
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<NormalizedRelayUrl, DrainFailure>()
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<HexKey>,
allLiveRelays: Set<NormalizedRelayUrl>,
bgScope: CoroutineScope,
) {
val missing = pubkeys.filter { relaysOf(it) == null }
if (missing.isEmpty()) return
suspend fun query(
authors: List<HexKey>,
relays: Set<NormalizedRelayUrl>,
) {
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<NormalizedRelayUrl>) {
val idsByAuthor = HashMap<HexKey, MutableList<HexKey>>()
for (ev in store.query<Event>(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<NormalizedRelayUrl, MutableSet<HexKey>>()
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<HexKey>,
backbone: Set<NormalizedRelayUrl>,
): Map<NormalizedRelayUrl, List<Filter>> {
val fallback = config.contentFallbackRelays
val perRelay = HashMap<NormalizedRelayUrl, MutableSet<HexKey>>()
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<Pair<List<HexKey>, Map<NormalizedRelayUrl, List<Filter>>>>(DRAIN_CONCURRENCY * 2)
val drainedOut = Channel<Triple<List<HexKey>, Set<NormalizedRelayUrl>, List<Pair<NormalizedRelayUrl, Event>>>>(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<NormalizedRelayUrl, DrainFailure>()
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<NormalizedRelayUrl, List<Filter>>,
deadOut: MutableMap<NormalizedRelayUrl, DrainFailure>?,
): List<Pair<NormalizedRelayUrl, Event>> {
if (filters.isEmpty()) return emptyList()
val eventChannel = Channel<Pair<NormalizedRelayUrl, Event>>(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<Pair<NormalizedRelayUrl, List<Filter>>>()
for ((relay, relayFilters) in filters) {
var group = ArrayList<Filter>()
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<NormalizedRelayUrl, DrainFailure>()
val timedOut = ConcurrentSet<NormalizedRelayUrl>()
val collected = mutableListOf<Pair<NormalizedRelayUrl, Event>>()
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<String>()
val groupListener =
object : SubscriptionListener {
override fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
forFilters: List<Filter>?,
) {
eventChannel.trySend(relay to event)
}
override fun onEose(
relay: NormalizedRelayUrl,
forFilters: List<Filter>?,
) {
done.complete("eose")
}
override fun onClosed(
message: String,
relay: NormalizedRelayUrl,
forFilters: List<Filter>?,
) {
done.complete("closed:$message")
}
override fun onCannotConnect(
relay: NormalizedRelayUrl,
message: String,
forFilters: List<Filter>?,
) {
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<Event>(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<Event>(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)
}
}
@@ -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<Int> = listOf(20, 10),
private val rateLadder: List<Long> = listOf(250L, 500L, 1000L, 2000L),
) : RelayConnectionListener {
private val gates = ConcurrentHashMap<NormalizedRelayUrl, Gate>()
private val gates = ConcurrentMap<NormalizedRelayUrl, Gate>()
// Concurrency-cap demotions per relay (== index+1 into subLadder). Capped at
// subLadder.size: past the floor we stop demoting.
private val subDemotions = ConcurrentHashMap<NormalizedRelayUrl, Int>()
private val subDemotions = ConcurrentMap<NormalizedRelayUrl, Int>()
// 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<NormalizedRelayUrl, Int>()
private val rateDelayMs = ConcurrentHashMap<NormalizedRelayUrl, Long>()
private val nextAllowedAtMs = ConcurrentHashMap<NormalizedRelayUrl, AtomicLong>()
private val rateSteps = ConcurrentMap<NormalizedRelayUrl, Int>()
private val rateDelayMs = ConcurrentMap<NormalizedRelayUrl, Long>()
private val nextAllowedAtMs = ConcurrentMap<NormalizedRelayUrl, AtomicLong>()
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<String, Any?> {
val cappedAt = sortedMapOf<Int, Int>()
for ((_, step) in subDemotions) {
val capCounts = HashMap<Int, Int>()
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<Long, Int>()
for ((_, step) in rateSteps) {
val cappedAt = capCounts.toList().sortedBy { it.first }.toMap()
val rateCounts = HashMap<Long, Int>()
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<CompletableDeferred<Unit>>()
@@ -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
}
}
}
@@ -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:<message>` 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
}
@@ -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<K : Any, V : Any>() {
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<K, V>
}
@@ -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<E : Any>() {
/** 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<E>
}
@@ -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<String, Int>()
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<String, Int>()
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<String, Int>()
// 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<String, Int>()
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<String>()
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<String>()
s.add("a")
val snap = s.snapshot()
s.add("b")
assertEquals(setOf("a"), snap)
assertEquals(2, s.size())
}
}
@@ -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<K : Any, V : Any> {
private val map = ConcurrentHashMap<K, V>()
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<K, V> = HashMap(map)
}
@@ -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<E : Any> {
private val set: MutableSet<E> = 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<E> = HashSet(set)
}
@@ -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<K : Any, V : Any> {
private val ref = AtomicReference(HashMap<K, V>())
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<K, V> = HashMap(ref.load())
}
@@ -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<E : Any> {
private val ref = AtomicReference(HashSet<E>())
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<E> = HashSet(ref.load())
}