Merge pull request #3511 from vitorpamplona/claude/graperank-sync-crawl-1n05im

graperank: relay reachability cache (NIP-66), aggregator kind:3 recovery, and outbox-discovery dedup
This commit is contained in:
Vitor Pamplona
2026-07-09 17:00:09 -04:00
committed by GitHub
11 changed files with 1267 additions and 123 deletions
@@ -71,6 +71,7 @@ import com.vitorpamplona.quartz.nip60Cashu.wallet.CashuWalletEvent
import com.vitorpamplona.quartz.nip61Nutzaps.info.NutzapInfoEvent
import com.vitorpamplona.quartz.nip61Nutzaps.nutzap.NutzapEvent
import com.vitorpamplona.quartz.nip65RelayList.AdvertisedRelayListEvent
import com.vitorpamplona.quartz.nip66RelayMonitor.reachability.RelayReachabilityStore
import com.vitorpamplona.quartz.nip87Ecash.recommendation.MintRecommendationEvent
import com.vitorpamplona.quartz.utils.SeenIds
import kotlinx.coroutines.CompletableDeferred
@@ -258,6 +259,21 @@ class Context(
private val storeDelegate: Lazy<IEventStore> = lazy { StoreFactory.open(dataDir) }
val store: IEventStore by storeDelegate
/**
* Shared relay-reachability cache (NIP-66 kind:30166 records in [store]), signed by
* the machine's dedicated monitor key — derived from the operator master, NOT the
* account (see [OperatorKeys.monitorKey]). The crawler and the WoT updater read its
* dead set to skip proven-dead relays and write their findings back, so liveness
* knowledge is shared across procedures and runs instead of rediscovered each time.
* Lazy so a run that never touches relays doesn't materialize the operator master.
*/
val reachability: RelayReachabilityStore by lazy {
RelayReachabilityStore(
store = store,
signer = NostrSignerInternal(dataDir.operatorKeys().monitorKey()),
)
}
/** Fully-wired manager. Call [prepare] once before use to load persisted state. */
val marmot: MarmotManager by lazy { MarmotManager(signer, mlsStore, messageStore, keyPackageStore) }
@@ -123,6 +123,24 @@ class OperatorKeys(
}
}
/**
* The machine's dedicated NIP-66 relay-monitor identity, derived once from the
* operator master (independent of any amy account). Unlike [serviceKey] this is
* NOT per-observer — the machine publishes relay-reachability (kind:30166) under a
* single, stable monitor pubkey, so a re-probe *replaces* the prior 30166 for a
* relay instead of orphaning it. Re-derivable from the one master seed alone.
*/
fun monitorKey(): KeyPair {
val master = masterPriv()
var counter = 0
while (true) {
val material = master + "$MONITOR_LABEL$counter".encodeToByteArray()
val kp = runCatching { KeyPair(privKey = sha256(material)) }.getOrNull()
if (kp?.privKey != null) return kp
counter++
}
}
private fun recordProvider(
observerHex: HexKey,
providerPubKey: HexKey,
@@ -152,5 +170,6 @@ class OperatorKeys(
private const val DIR_NAME = "operator"
private const val CONFIG_NAME = "operator.json"
private const val DERIVATION_LABEL = "graperank-provider:"
private const val MONITOR_LABEL = "relay-monitor:"
}
}
@@ -27,7 +27,7 @@ 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.GrapeRankCrawler
import com.vitorpamplona.quartz.experimental.graperank.GrapeRankParams
import com.vitorpamplona.quartz.experimental.graperank.GrapeRankPublisher
import com.vitorpamplona.quartz.experimental.graperank.GrapeRankUpdater
@@ -53,10 +53,17 @@ 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.CancellationException
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.asCoroutineDispatcher
import kotlinx.coroutines.async
import kotlinx.coroutines.awaitAll
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.withContext
import java.net.InetSocketAddress
import java.net.Socket
import java.net.URI
import java.util.concurrent.Executors
import kotlin.math.roundToInt
/**
@@ -78,13 +85,14 @@ import kotlin.math.roundToInt
*
* The crawl and the computation are separable, because the crawl persists every
* event it fetches to the store and the score is a pure function over it:
* - `amy graperank sync [OBSERVER]` — network only: crawl the reachable graph's
* kind 3/10000/1984/10002 into the local store. Idempotent and cumulative, so
* run it a few times to make sure everything is loaded. Scores nothing.
* - `amy graperank crawl [OBSERVER]` — network only: crawl the reachable graph's
* kind 3/10000/1984/10002 into the local store (aliased as the former `sync`).
* Idempotent and cumulative, so run it a few times to make sure everything is
* loaded. Scores nothing.
* - `amy graperank score [OBSERVER]` — local only: build the graph from the store
* and score (same as bare `--offline`). Instant and param-tunable; repeat with
* different `--rigor`/`--attenuation`/cutoffs without re-crawling.
* - bare `amy graperank [OBSERVER]` — the convenience combo: sync then score.
* - bare `amy graperank [OBSERVER]` — the convenience combo: crawl then score.
*
* Sub-verbs complete the NIP-85 provider experience — the discovery layer that
* lets clients find and consume those assertions:
@@ -105,6 +113,73 @@ object GrapeRankCommand {
"wss://eden.nostr.land",
).mapNotNull { RelayUrlNormalizer.normalizeOrNull(it) }.toSet()
// Network-wide aggregators that scrape and hold kind:3 for users whose own
// outbox lacks it. The crawler queries these for a straggler's CONTENT (kind:3),
// not just their kind:10002 relay list. Measured on observer 460c25e6, the distinct
// missing authors whose kind:3 each holds: kindpag.es 369, yabu 126, oxtr.dev 76,
// nos.lol 72, ditto 56, nostr1 29, momostr 11, mostr 3. So beyond the profile
// indexers (kindpag/purplepag/coracle/yabu/nostr1) and the ActivityPub bridges
// (ditto/momostr/mostr, which host bridged users' lists), two big general relays --
// nostr.oxtr.dev and nos.lol -- carry ~150 more that no indexer has.
private val CONTENT_AGGREGATOR_RELAYS: Set<NormalizedRelayUrl> =
DefaultIndexerRelayList +
listOf(
"wss://relay.ditto.pub",
"wss://relay.momostr.pink",
"wss://relay.mostr.pub",
"wss://nostr.oxtr.dev",
"wss://nos.lol",
).mapNotNull { RelayUrlNormalizer.normalizeOrNull(it) }.toSet()
private const val PROBE_TIMEOUT_MS = 2000
// The probe does BLOCKING DNS + TCP connect, and dead-domain DNS lookups can hang
// far past the connect timeout. On the shared Dispatchers.IO those hanging lookups
// starve the crawl's own IO — measured +462s on the finishing drain at hop-3. Run
// them on a dedicated, isolated daemon pool instead so the crawl's IO is untouched.
private val probeDispatcher =
Executors
.newFixedThreadPool(128) { r -> Thread(r, "relay-probe").apply { isDaemon = true } }
.asCoroutineDispatcher()
/**
* Cheap reachability pre-probe: a raw TCP connect (one round trip) with a tight
* timeout. Returns false only when the port won't even accept a socket — a dead
* dropper, refusal, or unroutable/onion/LAN host — which the crawler drops into
* deadHosts before the WS path pays its 7s connectTimeout. A busy-but-alive relay
* accepts the SYN instantly at the kernel level (its slowness is at the app layer),
* so it passes here and is left for the real WS attempt. Unparseable host → true,
* so an odd URL is never culled on a parse quirk — let the WS decide.
*/
private suspend fun tcpReachable(relay: NormalizedRelayUrl): Boolean =
withContext(probeDispatcher) {
val hostPort = relayHostPort(relay) ?: return@withContext true
try {
Socket().use { it.connect(InetSocketAddress(hostPort.first, hostPort.second), PROBE_TIMEOUT_MS) }
true
} catch (e: Exception) {
if (e is CancellationException) throw e
false
}
}
private fun relayHostPort(relay: NormalizedRelayUrl): Pair<String, Int>? =
try {
val uri = URI(relay.url)
val host = uri.host ?: return null
val port =
if (uri.port > 0) {
uri.port
} else if (relay.url.startsWith("wss://", ignoreCase = true)) {
443
} else {
80
}
host to port
} catch (e: Exception) {
null
}
suspend fun dispatch(
dataDir: DataDir,
tail: Array<String>,
@@ -115,7 +190,9 @@ object GrapeRankCommand {
"register" -> register(dataDir, tail.drop(1).toTypedArray())
"providers" -> providers(dataDir, tail.drop(1).toTypedArray())
"operator" -> operator(dataDir, tail.drop(1).toTypedArray())
"sync" -> sync(dataDir, tail.drop(1).toTypedArray())
// `sync` is the pre-rename name kept as a back-compat alias; `crawl` is
// canonical (disambiguates from negentropy `amy sync` / `graperank update`).
"crawl", "sync" -> crawl(dataDir, tail.drop(1).toTypedArray())
"update" -> update(dataDir, tail.drop(1).toTypedArray())
"score" -> run(dataDir, tail.drop(1).toTypedArray(), forceOffline = true)
else -> run(dataDir, tail)
@@ -174,12 +251,13 @@ object GrapeRankCommand {
// 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
var crawlStats: GrapeRankCrawler.Stats? = null
if (!offline) {
val stats = newCrawler(ctx, args).crawl(observer, builder)
crawlStats = stats
contactListsFed = stats.contactListsFed
flushReachability(ctx, args, stats)
reportRelayFeedback(ctx)
} else {
// Offline: stream contact lists from the local store into the graph.
@@ -347,7 +425,7 @@ object GrapeRankCommand {
/**
* Configure the outbox-model crawler from the crawl flags on [args] plus the
* account's relay policy. Shared by the bare command and `graperank sync`.
* account's relay policy. Shared by the bare command and `graperank crawl`.
* Relay policy — where a stranger's kind:10002 is found (index/discovery
* aggregators + general defaults) and best-effort general relays that might
* hold content when an outbox is unknown — lives in app code, so the quartz
@@ -356,18 +434,28 @@ object GrapeRankCommand {
private suspend fun newCrawler(
ctx: Context,
args: Args,
): GrapeRankDataCrawler {
): GrapeRankCrawler {
val discoveryRelays =
ctx.bootstrapRelays() + Constants.eventFinderRelays + DefaultIndexerRelayList + EXTRA_DISCOVERY_RELAYS
val contentFallback = ctx.bootstrapRelays() + Constants.eventFinderRelays
return GrapeRankDataCrawler(
// Aggregator kind:3 recovery for stragglers is on by default; --no-aggregators
// disables it for A/B comparison.
val aggregators = if (args.bool("no-aggregators")) emptySet() else CONTENT_AGGREGATOR_RELAYS
// Seed the crawl with relays a prior run/monitor proved dead within the cache's
// TTL, so the WS path never re-pays their connect timeouts (--no-reachability-cache
// to skip). The crawl's own final live/dead set is flushed back by the caller.
val knownDead =
if (args.bool("no-reachability-cache")) emptySet() else ctx.reachability.snapshot().dead
return GrapeRankCrawler(
client = ctx.client,
store = ctx.store,
limiter = ctx.relayLimiter,
config =
GrapeRankDataCrawler.Config(
GrapeRankCrawler.Config(
relayListDiscoveryRelays = discoveryRelays,
knownDeadRelays = knownDead,
contentFallbackRelays = contentFallback,
contentAggregatorRelays = aggregators,
maxRounds = args.intFlag("max-rounds", Int.MAX_VALUE),
maxHops = args.intFlag("max-hops", Int.MAX_VALUE),
timeoutMs = args.longFlag("timeout", 10L) * 1000,
@@ -375,6 +463,11 @@ object GrapeRankCommand {
diagnose = args.bool("diagnose"),
insertBatchSize = args.intFlag("insert-batch", 500),
drainConcurrency = args.intFlag("drain-concurrency", 24),
timeoutEvictStrikes = args.intFlag("timeout-evict", 3),
// Cheap TCP reachability pre-probe (--no-probe to disable). No Tor
// transport here, so .onion relays are skipped on sight.
reachabilityProbe = if (args.bool("no-probe")) null else ::tcpReachable,
torEnabled = false,
// shedDeadDiscovery / shardRotations keep their benchmarked-best
// Config defaults.
),
@@ -393,12 +486,33 @@ object GrapeRankCommand {
}
/**
* `amy graperank sync [OBSERVER]` — network-only WoT data sync. Crawls the
* reachable follow/mute/report graph into the local store (kind 3/10000/1984/
* 10002) and reports what it loaded, WITHOUT scoring. Idempotent + cumulative:
* run it a few times to make sure everything is loaded, then `graperank score`.
* Flush the crawl's final live/dead relay verdicts into the shared reachability
* cache (NIP-66 kind:30166) so the next crawl and the WoT updater start warm and
* skip proven-dead relays. Best-effort and behind `--no-reachability-cache`: a
* cache write must never fail the crawl it is summarizing.
*/
private suspend fun sync(
private suspend fun flushReachability(
ctx: Context,
args: Args,
stats: GrapeRankCrawler.Stats,
) {
if (args.bool("no-reachability-cache")) return
runCatching {
ctx.reachability.record(reachable = stats.liveRelays, dead = stats.deadRelays)
System.err.println(
"[graperank] reachability cache: recorded ${stats.liveRelays.size} live, ${stats.deadRelays.size} dead",
)
}.onFailure { System.err.println("[graperank] reachability cache flush failed: ${it.message}") }
}
/**
* `amy graperank crawl [OBSERVER]` — network-only WoT data crawl (aliased as the
* former `sync`). Crawls the reachable follow/mute/report graph into the local
* store (kind 3/10000/1984/10002) and reports what it loaded, WITHOUT scoring.
* Idempotent + cumulative: run it a few times to make sure everything is loaded,
* then `graperank score`.
*/
private suspend fun crawl(
dataDir: DataDir,
rest: Array<String>,
): Int {
@@ -410,6 +524,7 @@ object GrapeRankCommand {
// Persist-only crawl: no in-memory graph (null builder); every event
// still lands in the store for a later `score`.
val stats = newCrawler(ctx, args).crawl(observer, null)
flushReachability(ctx, args, stats)
reportRelayFeedback(ctx)
Output.emit(
linkedMapOf<String, Any?>(
@@ -466,6 +581,13 @@ object GrapeRankCommand {
Context.openOrAnonymous(dataDir).use { ctx ->
ctx.prepare()
// Skip relays a crawl/monitor proved dead within the cache's TTL — a dead
// relay cannot serve its authors, so reconciling it only burns a timeout.
// Live author-advertised relays are always synced (--no-reachability-cache
// to reconcile every relay regardless).
val knownDead =
if (args.bool("no-reachability-cache")) emptySet() else ctx.reachability.snapshot().dead
val updater =
GrapeRankUpdater(
client = ctx.client,
@@ -479,6 +601,7 @@ object GrapeRankCommand {
authorChunk = args.intFlag("author-chunk", 500),
minAuthors = args.intFlag("min-authors", 1),
idleTimeoutMs = args.longFlag("timeout", 30L) * 1000,
knownDead = knownDead,
),
log = { System.err.println(it) },
)
@@ -491,7 +614,7 @@ object GrapeRankCommand {
"relay_lists_in_store" to result.relayListsInStore,
"authors_with_outbox" to result.authorsWithOutbox,
"relays" to 0,
"note" to "no kind:10002 write relays in the local store — run `graperank sync` first",
"note" to "no kind:10002 write relays in the local store — run `graperank crawl` first",
),
)
return 0
@@ -27,10 +27,12 @@ import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.AdaptiveRelayLimiter
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.DrainFailure
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.classifyDrainFailure
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAllPages
import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener
import com.vitorpamplona.quartz.nip01Core.relay.client.single.newSubId
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import com.vitorpamplona.quartz.nip01Core.store.IEventStore
import com.vitorpamplona.quartz.nip02FollowList.ContactListEvent
import com.vitorpamplona.quartz.nip09Deletions.DeletionEvent
@@ -48,10 +50,13 @@ import kotlinx.coroutines.awaitAll
import kotlinx.coroutines.cancel
import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.currentCoroutineContext
import kotlinx.coroutines.delay
import kotlinx.coroutines.isActive
import kotlinx.coroutines.joinAll
import kotlinx.coroutines.launch
import kotlinx.coroutines.selects.select
import kotlinx.coroutines.sync.Semaphore
import kotlinx.coroutines.withTimeoutOrNull
import kotlin.concurrent.atomics.AtomicLong
import kotlin.concurrent.atomics.ExperimentalAtomicApi
@@ -86,7 +91,7 @@ import kotlin.time.TimeSource
* is emitted through [log]; a headless caller routes it to stderr, a UI ignores it.
*/
@OptIn(ExperimentalAtomicApi::class)
class GrapeRankDataCrawler(
class GrapeRankCrawler(
private val client: NostrClient,
private val store: IEventStore,
private val limiter: AdaptiveRelayLimiter,
@@ -111,6 +116,15 @@ class GrapeRankDataCrawler(
* 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 contentAggregatorRelays index/aggregator relays that hold a network-wide
* copy of kind:3, mined for stragglers by [recoverStragglersFromAggregators]
* after the crawl converges. The outbox model asks "where does this user write?"
* — but a large tail of users have no kind:3 on their own advertised outbox (it's
* dead, or they never published one there), while a network-wide aggregator
* (kindpag.es, …) scraped and holds it. Those aggregators are queried only for
* kind:10002 in [ensureRelayLists]; the dedicated recovery pass asks them for the
* kind:3 itself — patiently and kind:3-only, since a multi-kind filter makes the
* big aggregators time out. Empty disables the pass.
* @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 the FAST per-drain timeout that gates a round's progression.
@@ -134,10 +148,18 @@ class GrapeRankDataCrawler(
* moderate: a higher global fan-out re-floods busy hubs faster than demotion
* catches up (an A/B at 64 ran ~2x slower with more dead relays), so 24 is the
* validated default and raising it is a probe, not a speedup.
* @param timeoutEvictStrikes evict a relay after this many drains that timed out
* (connect timeout or park idle-cut) having delivered NOTHING. Unlike
* [classifyDrainFailure] — which never marks a timeout dead, since one slow
* answer shouldn't drop a relay — this catches the connect-but-silent / dead
* endpoints that are otherwise re-tried through every straggler's outbox for the
* rest of the crawl. A clean EOSE or any delivered event clears a relay's count,
* so only never-productive relays are evicted. `<= 0` disables it.
*/
class Config(
val relayListDiscoveryRelays: Set<NormalizedRelayUrl>,
val contentFallbackRelays: Set<NormalizedRelayUrl>,
val contentAggregatorRelays: Set<NormalizedRelayUrl> = emptySet(),
val maxRounds: Int = Int.MAX_VALUE,
val maxHops: Int = Int.MAX_VALUE,
val timeoutMs: Long = 10_000,
@@ -145,6 +167,7 @@ class GrapeRankDataCrawler(
val diagnose: Boolean = false,
val insertBatchSize: Int = 500,
val drainConcurrency: Int = 24,
val timeoutEvictStrikes: Int = 3,
/**
* Also skip proven-dead relays in the kind:10002 discovery sweep
* ([ensureRelayLists]). Without it the discovery/backbone set is queried
@@ -158,6 +181,33 @@ class GrapeRankDataCrawler(
* 6-rotation wall cost with no completeness loss (1 clears too little).
*/
val shardRotations: Int = 2,
/**
* Optional cheap reachability pre-probe. Given a relay, returns false if it
* is definitely unreachable from here — a raw TCP connect (one round trip)
* that failed fast. A background culler runs it over the cold tail of learned
* relays and drops the unreachable ones into [deadHosts] BEFORE the expensive
* WS path pays the full 7s connectTimeout on them. It only ever marks dead
* and defers to the WS verdict: a host already proven live/dead is skipped.
* A tight TCP timeout is safe where a tight WS timeout is not — a busy-but-
* alive relay accepts the SYN instantly (kernel-level) and only stalls at the
* app layer, so TCP-reachability separates "unreachable" from "slow". Null
* disables pre-probing.
*/
val reachabilityProbe: (suspend (NormalizedRelayUrl) -> Boolean)? = null,
/**
* Whether this client can reach .onion relays (has a Tor transport). When
* false, every .onion relay is unreachable and [isDead] skips it on sight —
* no socket, no wasted connect attempt.
*/
val torEnabled: Boolean = false,
/**
* Relays a prior run (or another monitor) proved unreachable within the
* reachability cache's TTL — seeded into [deadRelays] before the crawl starts
* so we don't re-pay their connect timeouts. Per-URL (not per-authority): a
* TTL'd "skip for now", re-probed once the record ages out, so it never
* permanently ignores an author's advertised home. See RelayReachabilityStore.
*/
val knownDeadRelays: Set<NormalizedRelayUrl> = emptySet(),
)
/** What the crawl fetched — the counters the caller reports and the graph is built from. */
@@ -180,6 +230,17 @@ class GrapeRankDataCrawler(
val insertMs: Long,
/** Verified events handed to the store (duplicates included — the write path dedups). */
val eventsStored: Long,
/**
* Relays this run actually OBSERVED, for the caller to flush into the
* reachability cache (kind:30166): [deadRelays] = relays newly proven
* unreachable this run (a connect-establishment failure we paid), [liveRelays]
* = relays that served ≥1 event. Seeded known-dead relays (skipped, never
* dialed) are deliberately EXCLUDED from [deadRelays] — re-writing them would
* refresh their TTL without a re-probe and blacklist a recovered relay forever;
* their original record must age out so the next run re-probes them.
*/
val deadRelays: Set<NormalizedRelayUrl>,
val liveRelays: Set<NormalizedRelayUrl>,
)
/**
@@ -201,13 +262,13 @@ class GrapeRankDataCrawler(
/**
* Holds all per-crawl mutable state. Graph state (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. The frontier IS [hopOf]'s key set: a user
* is "discovered" iff it has a hop stamp. Only the state genuinely shared across
* the producer / consumer / drain-worker coroutines is concurrent: relayHints,
* attempts, deadRelays, relayStrikes.
* relaysContacted) is single-writer by construction — Phase A and the Phase-B
* consumer never run concurrently — so those stay plain collections. The frontier
* IS [hopOf]'s key set: a user is "discovered" iff it has a hop stamp. State
* genuinely shared across the producer / consumer / drain-worker coroutines is
* concurrent: relayHints, attempts, deadRelays. [writeRelayFreq] and [liveRelays]
* are also concurrent because the background reachability culler reads them (to
* find candidates and skip already-live authorities) while the crawl writes them.
*/
private inner class CrawlRun(
val observer: HexKey,
@@ -217,9 +278,25 @@ class GrapeRankDataCrawler(
// is the discovered frontier — no separate `discovered` set to keep in sync.
val hopOf = hashMapOf(observer to 0)
val done = hashSetOf<HexKey>()
// Users we've already run the STATIC-relay outbox-discovery sweep for (the
// indexer set in [ensureRelayLists]). That relay set never changes, so re-asking
// it for the same never-had-a-10002 user each round it recirculates is pure waste
// — the round-8 profile showed this as ~144k slow kind:10002 drains.
// Single-writer: only the round loop's [ensureRelayLists] touches it.
val relayListDiscoverySwept = hashSetOf<HexKey>()
// Relays the WIDE (every-live-relay) recovery pass in [ensureRelayLists] has
// already asked. Unlike the discovery set the wide net GROWS as the crawl learns
// relays, so gating that pass on the swept-user set would never re-ask an old
// straggler for a home relay discovered after its first sweep. Tracking the asked
// relay set instead lets a late-appearing relay still surface an old straggler's
// 10002 while never asking the same (user, relay) pair twice.
// Single-writer: only the round loop's [ensureRelayLists] touches it.
val wideRelaysSwept = hashSetOf<NormalizedRelayUrl>()
val relaysContacted = hashSetOf<NormalizedRelayUrl>()
val writeRelayFreq = HashMap<NormalizedRelayUrl, Int>()
val liveRelays = hashSetOf<NormalizedRelayUrl>()
val writeRelayFreq = ConcurrentMap<NormalizedRelayUrl, Int>()
val liveRelays = ConcurrentSet<NormalizedRelayUrl>()
// Per-relay outcome/latency/yield accounting, written from every drain unit
// (fast + parked) across every round. Dumped at crawl end; the raw signal a
@@ -229,8 +306,39 @@ class GrapeRankDataCrawler(
// 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>()
// Seeded from the reachability cache (relays proven dead within its TTL) so the
// WS path never re-pays their connect timeouts; the crawl still adds/removes
// more as it goes and flushes the union back at the end.
val deadRelays = ConcurrentSet<NormalizedRelayUrl>().apply { config.knownDeadRelays.forEach { add(it) } }
// Unproductive-TIMEOUT strikes, keyed by relay AUTHORITY (host[:port]), not the
// full URL. [classifyDrainFailure] treats a READ timeout or a park idle-cut — the
// relay answered the handshake but is slow — as "busy, retry" and never dead,
// because one slow answer shouldn't evict a relay. But in a crawl the same
// unresponsive server is
// routed through every straggler's outbox, every round, each visit burning the
// full timeout + park window for zero data. Keying by authority is what defeats
// the outbox-model's per-user path fragmentation: a paid/dead host like
// `filter.nostr.wine` is advertised as hundreds of distinct per-user URLs
// (`filter.nostr.wine/npubA?broadcast=true`, …), so a per-URL counter never
// reaches the threshold on any single one — but they are one server, and it
// times out on all of them. We strike the authority and, past
// [Config.timeoutEvictStrikes], mark it dead in [deadHosts] so every URL under
// it is skipped. A clean EOSE or any delivered event records the authority in
// [producedHosts] (see [clearTimeoutStrikes]), and [isDead] treats an authority
// as dead ONLY while it is in [deadHosts] AND NOT in [producedHosts] — so a host
// that ever produces is never evicted, even if concurrent strikes from the
// 24-worker fan-out raced it into [deadHosts] at the same instant it EOSE'd.
// Only the connect-but-silent / dead-endpoint class stays evicted.
val deadTimeoutStrikes = ConcurrentMap<String, Int>()
val deadHosts = ConcurrentSet<String>()
// Authorities that ever produced (EOSE or a delivered event) this run. Membership
// here overrides [deadHosts] in [isDead], making the "ever produces ⇒ never
// evicted" invariant race-free: the strike path can lose to a clear and still add
// to [deadHosts], but the gate consults this set and lets the proven host run.
val producedHosts = ConcurrentSet<String>()
// Crawl-wide dedup of event ids, shared across all concurrent drains and
// every round. The outbox model mirrors the SAME event (especially kind:10002
@@ -262,6 +370,24 @@ class GrapeRankDataCrawler(
// it finishes (or its park window elapses).
val parkedInFlight = AtomicLong(0)
// ── Saturation / latency instrumentation (diagnose only) ─────────────────────
// Are we resource-bound or waiting-on-relays? These answer it without a profiler.
// activeWorkers: Phase-B drain workers busy right now (vs drainConcurrency) — if
// rarely full, adding workers won't help; the producer/relays are the limit.
// throttled: drains that got a relay rate-limit (429/too-many-*) — the EXTERNAL
// ceiling; if it climbs when we push harder, more concurrency backfires.
// The latency sums split a drain's wall time into time-to-first-event vs the
// EOSE-wait AFTER the relay's last event (pure waiting on a done-but-slow relay)
// — a large eose-wait fraction is the case for a shorter/adaptive fast window.
// burnedFastWindow: drains that blew timeoutMs and had to park.
val activeWorkers = AtomicLong(0)
val throttled = AtomicLong(0)
val firstEventSumMs = AtomicLong(0)
val eoseWaitSumMs = AtomicLong(0)
val drainWallSumMs = AtomicLong(0)
val drainSamples = AtomicLong(0)
val burnedFastWindow = AtomicLong(0)
// Background scope owning the parked subscriptions (and Tier-2 relay-list
// sweeps). Set in [run]; cancelled once the crawl converges.
var bgScope: CoroutineScope? = null
@@ -287,32 +413,132 @@ class GrapeRankDataCrawler(
var progConverging = false
/**
* 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.
* A relay [classifyDrainFailure] flagged [DrainFailure.DEAD] won't serve us
* this run (bad domain, TLS misconfig, dead/gated HTTP code, refused/reset,
* connect that never opened), so it is dropped on the first strike. Read
* timeouts and alive 429 rate-limits 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)
}
}
for ((r, _) in failed) deadRelays.add(r)
}
/**
* A relay's drain unit timed out (connect timeout or park idle-cut) having
* delivered nothing. Count the strike against its AUTHORITY and, once it
* reaches [Config.timeoutEvictStrikes], give up on the whole host — a
* connect-but-silent or dead endpoint that would otherwise be re-tried through
* every straggler's outbox for the rest of the crawl. Disabled when the
* threshold is <= 0.
*/
fun strikeUnproductiveTimeout(relay: NormalizedRelayUrl) {
val limit = config.timeoutEvictStrikes
if (limit <= 0) return
val authority = authorityOf(relay.url)
if (authority in producedHosts) return // proven productive — never evict on timeouts
if (authority in deadHosts) return
if (deadTimeoutStrikes.merge(authority, 1) { a, b -> a + b } >= limit) deadHosts.add(authority)
}
/**
* A relay just proved its host can produce — a clean EOSE or an actual event —
* so record the authority in [producedHosts] (permanently protecting it from
* timeout eviction for the rest of the run) and wipe any timeout strikes it
* accrued. Prevents an occasionally-slow but useful host (a busy backbone hub, or
* a multi-path relay where some paths are slow) from accumulating its way to
* eviction across a long crawl — and, via the [isDead] gate, un-evicts one that a
* concurrent strike already pushed into [deadHosts] at the same instant.
*/
fun clearTimeoutStrikes(relay: NormalizedRelayUrl) {
val authority = authorityOf(relay.url)
producedHosts.add(authority)
// ConcurrentMap exposes no remove; reset the count to 0 atomically (0 is
// below any positive eviction threshold, so it reads as "unstruck"). Guard
// on a prior entry so we don't insert a 0 for every host that ever answers.
if (deadTimeoutStrikes[authority] != null) deadTimeoutStrikes.merge(authority, 0) { _, _ -> 0 }
}
/**
* A relay is out of the routing pool if it hard/transient-failed (per-URL
* [deadRelays]) or its whole authority was timeout-evicted ([deadHosts]) and has
* not since proven productive ([producedHosts] wins, so a slow-but-live host is
* never permanently evicted); a .onion relay is dead on sight unless we have a
* Tor transport, since every connect to it would only hang and fail.
*/
fun isDead(relay: NormalizedRelayUrl): Boolean {
val authority = authorityOf(relay.url)
return relay in deadRelays ||
(authority in deadHosts && authority !in producedHosts) ||
(!config.torEnabled && RelayUrlNormalizer.isOnion(relay.url))
}
/** The busiest live relays we've learned, excluding the dead ones. */
fun topLiveRelays(cap: Int): List<NormalizedRelayUrl> =
writeRelayFreq.entries
writeRelayFreq
.snapshot()
.entries
.asSequence()
.filter { it.key in liveRelays && it.key !in deadRelays }
.filter { it.key in liveRelays && !isDead(it.key) }
.sortedByDescending { it.value }
.take(cap)
.map { it.key }
.toList()
/**
* Background reachability culler. Cheaply TCP-probes the relays we've learned —
* COLD TAIL FIRST — and drops the unreachable ones into [deadHosts] so the WS
* path never pays the 7s connectTimeout on a dead host. It only ever marks dead
* and probes each authority once. Any host the WS path already resolved is
* skipped: dead ones via [isDead], and hosts already proven LIVE ([liveRelays])
* are filtered out up front so we never waste a probe — or a needless TCP hit —
* on a working relay we depend on. Combined with the cold-tail ordering, the
* probe stays off the hot relays the crawl is actively dialing, and the WS
* verdict always wins ("if the websocket gets there first, let it run"). Runs on
* [bgScope] until the crawl cancels it.
*/
private suspend fun cullUnreachable(probe: suspend (NormalizedRelayUrl) -> Boolean) {
val probed = HashSet<String>() // authorities; only ever touched by this coroutine's loop
val gate = Semaphore(PROBE_CONCURRENCY)
while (currentCoroutineContext().isActive) {
// Never probe a host the WS path already proved live — wasted work and
// a needless TCP hit on the hot relays we depend on.
val liveAuthorities = liveRelays.snapshot().mapTo(HashSet()) { authorityOf(it.url) }
// Least-written relays are the niche/dead long tail the WS path reaches
// last — probing them first buys the most head start with the least
// contention against the busy relays already being connected.
val batch =
writeRelayFreq
.snapshot()
.entries
.asSequence()
.filter {
val authority = authorityOf(it.key.url)
authority !in probed && authority !in liveAuthorities && !isDead(it.key)
}.sortedBy { it.value }
.map { it.key }
.toList()
if (batch.isEmpty()) {
delay(PROBE_IDLE_MS)
continue
}
coroutineScope {
for (relay in batch) {
val authority = authorityOf(relay.url)
if (!probed.add(authority)) continue
gate.acquire()
launch {
try {
// Re-check: the WS path may have resolved it while queued.
if (!isDead(relay) && !probe(relay)) deadHosts.add(authority)
} finally {
gate.release()
}
}
}
}
}
}
/**
* 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.
@@ -361,6 +587,26 @@ class GrapeRankDataCrawler(
return got
}
/**
* Mark done any still-pending user whose kind:3 is already in the store — a
* previous round's [ensureRelayLists] co-fetch, a late parked delivery, or a
* prior run's data — folding it into the graph so Phase B never spends an outbox
* drain re-pulling a contact list we already hold. Single-writer: called only
* from the round loop at Phase-A time, before the drain workers start. Returns
* the count newly fed.
*/
suspend fun harvestFromStore(authors: Collection<HexKey>): Int {
var got = 0
for (pk in authors) {
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
@@ -451,18 +697,16 @@ class GrapeRankDataCrawler(
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>,
kinds: List<Int>,
) {
if (relays.isEmpty() || authors.isEmpty()) return
val filters =
relays.associateWith {
authors.chunked(AUTHORS_PER_FILTER).map { chunk ->
Filter(kinds = listOf(AdvertisedRelayListEvent.KIND), authors = chunk)
Filter(kinds = kinds, authors = chunk)
}
}
drainGated(filters, null)
@@ -474,12 +718,40 @@ class GrapeRankDataCrawler(
} else {
config.relayListDiscoveryRelays
}
query(missing, discovery)
val stillMissing = missing.filter { relaysOf(it) == null }
// First-time DISCOVERY sweep on the static indexer set: a user only needs it
// once (re-asking a static set can't find a 10002 we already missed). Co-fetch
// the contact list in the same REQ — an indexer holding a user's 10002 often
// holds their kind:3, a cheap byproduct of a round-trip we already pay.
val freshlyMissing = pubkeys.filter { relaysOf(it) == null && it !in relayListDiscoverySwept }
val freshlyMissingSet = freshlyMissing.toHashSet()
relayListDiscoverySwept.addAll(freshlyMissing)
query(freshlyMissing, discovery, listOf(AdvertisedRelayListEvent.KIND, ContactListEvent.KIND))
// WIDE recovery sweep for kind:10002 ONLY (co-fetching kind:3 across thousands
// of relays inflates this fire-and-forget sweep, which the finishing drain then
// waits on — measured +300s at hop-3). The wide net grows every round, so it is
// gated on the asked-RELAY set, not the swept-user set:
// - a user first swept this round is asked the whole current wide net;
// - a user swept earlier and still missing is asked ONLY the relays that
// appeared since — so its home relay, discovered late, still surfaces.
// No (user, relay) pair is asked twice; every straggler eventually sees every
// live relay. The added olderStillMissing×newWide work self-limits: newWide
// shrinks toward zero as the relay universe is exhausted.
val wide = allLiveRelays - discovery
if (stillMissing.isNotEmpty() && wide.isNotEmpty()) {
bgScope.launch { query(stillMissing, wide) }
val newWide = wide - wideRelaysSwept
wideRelaysSwept.addAll(wide)
val freshStillMissing = freshlyMissing.filter { relaysOf(it) == null }
val olderStillMissing = pubkeys.filter { it !in freshlyMissingSet && relaysOf(it) == null }
val hasWork =
(freshStillMissing.isNotEmpty() && wide.isNotEmpty()) ||
(olderStillMissing.isNotEmpty() && newWide.isNotEmpty())
if (hasWork) {
bgScope.launch {
query(freshStillMissing, wide, listOf(AdvertisedRelayListEvent.KIND))
query(olderStillMissing, newWide, listOf(AdvertisedRelayListEvent.KIND))
}
}
}
@@ -502,7 +774,7 @@ class GrapeRankDataCrawler(
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)
for (relay in write) if (!isDead(relay)) perRelayAuthors.getOrPut(relay) { HashSet() }.add(author)
}
if (perRelayAuthors.isEmpty()) return
@@ -531,6 +803,12 @@ class GrapeRankDataCrawler(
* [backbone] — the known-good relays other people write to;
* - no outbox at all: harvested hints + backbone + the general fallback.
*
* The content aggregators are deliberately NOT mixed in here: this path's
* multi-kind [FETCH_KINDS] query loses their kind:3 to their per-REQ result
* cap (a big indexer fills the response with the abundant kind:10002 and
* returns no kind:3), so recovering from them is done separately — kind:3-only,
* once and patiently — in [recoverStragglersFromAggregators].
*
* Also tallies each user's write relays into [writeRelayFreq] so the
* backbone can be learned from the crawl. Authors are chunked per relay.
*/
@@ -543,7 +821,7 @@ class GrapeRankDataCrawler(
for (pk in pubkeys) {
val write = relaysOf(pk)?.writeRelaysNorm()?.takeIf { it.isNotEmpty() }
write?.forEach { writeRelayFreq[it] = (writeRelayFreq[it] ?: 0) + 1 }
write?.forEach { writeRelayFreq.merge(it, 1) { a, b -> a + b } }
val relays =
when {
write == null -> relayHints[pk]?.snapshot().orEmpty() + backbone + fallback
@@ -555,7 +833,7 @@ class GrapeRankDataCrawler(
// (re-querying them for this user is guaranteed-empty waste).
val emptied = askedEmpty[pk]
for (relay in relays) {
if (relay in deadRelays) continue
if (isDead(relay)) continue
if (emptied != null && relay in emptied) continue
perRelay.getOrPut(relay) { HashSet() }.add(pk)
}
@@ -568,6 +846,87 @@ class GrapeRankDataCrawler(
}
}
/**
* Final patient pass for the stragglers the outbox model couldn't resolve.
* A large tail of reachable users have no kind:3 on their own advertised
* outbox — it's dead, or they never published one there — while a
* network-wide aggregator ([Config.contentAggregatorRelays], e.g.
* kindpag.es) scraped and holds it. Mixing those aggregators into the
* competitive Phase-B fan-out doesn't work: there they'd be asked for the
* multi-kind [FETCH_KINDS] filter (which times them out) and would race
* thousands of outbox sockets, getting cut before a big aggregator finishes.
* So once the frontier is drained we ask the aggregators for the remaining
* stragglers' kind:3 ALONE: a handful of relays drained kind:3-only with the
* patient park window, not competing with the fan-out. Recovered contact
* lists are folded into the graph and persisted for a later `score`.
*/
private suspend fun recoverStragglersFromAggregators() {
// Query EVERY configured aggregator, even ones the main crawl evicted.
// During the competitive crawl an indexer like user.kindpag.es is only ever
// asked for kind:10002 in bulk and kind:[3,10000,1984,10002] one author at a
// time; the latter parks and times out (60–80s each), striking the host until
// it's timeout-evicted (isDead). It is never asked for a clean bulk kind:3 —
// the one thing it actually serves fast (≈19 lists per 300 authors in a few
// seconds). This deliberate patient pass IS that clean query, so eviction from
// the fan-out must not disqualify it here. [drainGated] subscribes to whatever
// filter map we hand it (it does not re-check isDead), and a genuinely dead
// endpoint just costs one shared park window since the units run concurrently.
val aggregators = config.contentAggregatorRelays.toHashSet()
if (aggregators.isEmpty()) return
// Wipe any timeout strikes the fan-out accrued so a partially-struck host
// starts this pass clean and a fast EOSE here keeps it healthy.
for (agg in aggregators) clearTimeoutStrikes(agg)
// Stragglers = crawled users we still have no kind:3 for. Most are already
// in `done` (their outbox attempts were exhausted), which is exactly why
// [harvest]/[ingestLate] can't be reused — they skip `done` users — so we
// fold these directly.
val stragglers = hopOf.keys.filterTo(HashSet()) { (hopOf[it] ?: 0) < config.maxHops && contactsOf(it) == null }
if (stragglers.isEmpty()) return
val before = contactListsFed
log("[graperank] aggregator recovery: ${stragglers.size} stragglers via ${aggregators.size} aggregators")
// Build the query against the full straggler set BEFORE any folding (the
// filter lists are materialized here, so later mutation of `stragglers` is
// safe). Ask ONLY for kind:3 — the contact list we're missing. A multi-kind
// filter is useless against the big indexers: user.kindpag.es caps its
// response at ~100 events per REQ (it ignores our limit), so a
// kinds=[3,10000,1984,10002] query comes back 100× kind:10002 and 0×
// kind:3 — the abundant relay lists crowd the contact lists out entirely.
// Asked for kind:3 alone it returns them in a few seconds. Their kind:10002
// is already fetched in bulk by [ensureRelayLists]; mutes/reports still come
// from the outbox model. The aggregator's job here is only the lists.
val filters =
aggregators.associateWith {
stragglers.chunked(AUTHORS_PER_FILTER).map { chunk -> Filter(kinds = listOf(ContactListEvent.KIND), authors = chunk) }
}
// Fold one delivered contact list per straggler, exactly once.
suspend fun foldAgg(events: List<Pair<NormalizedRelayUrl, Event>>) {
for ((relay, ev) in events) {
liveRelays.add(relay)
if (ev !is ContactListEvent) continue
val pk = ev.pubKey
if (pk !in stragglers) continue
val contacts = contactsOf(pk) ?: continue
stragglers.remove(pk)
done += pk
ingest(pk, contacts)
}
}
relaysContacted += aggregators
// Fast deliveries fold immediately; a slow aggregator parks and its late
// kind:3 arrives on [lateHarvest], which we drain until the parked units
// finish — so a big aggregator that can't answer within the fast window is
// still fully harvested here instead of being abandoned.
foldAgg(drainGated(filters, null))
while (parkedInFlight.load() > 0L) {
withTimeoutOrNull(PARK_POLL_MS) { lateHarvest.receive() }?.let { foldAgg(listOf(it)) }
}
while (true) foldAgg(listOf(lateHarvest.tryReceive().getOrNull() ?: break))
log("[graperank] aggregator recovery: +${contactListsFed - before} contact lists")
}
/**
* Dedup (crawl-wide [seenIds]), verify, and group-commit a unit's events,
* returning the newly-stored ones tagged by relay. Safe to call concurrently
@@ -599,7 +958,7 @@ class GrapeRankDataCrawler(
val ok = event.verify()
verifyNanos.addAndFetch(vMark.elapsedNow().inWholeNanoseconds)
if (!ok) {
Log.w("GrapeRankDataCrawler") { "dropped event ${event.id.take(8)} kind=${event.kind} — bad signature" }
Log.w("GrapeRankCrawler") { "dropped event ${event.id.take(8)} kind=${event.kind} — bad signature" }
continue
}
if (!seenIds.add(event.id)) continue // lost the race to a mirror; it stores it
@@ -701,15 +1060,57 @@ class GrapeRankDataCrawler(
val pct = (100L * roundDone / progTarget).coerceIn(0, 100)
val remaining = (progTarget - roundDone).coerceAtLeast(0)
val eta = if (rate > 0) etaFmt(remaining / rate) else "…"
// Saturation tail: workers busy / cap, and rate-limit hits so far — is
// the pool full (raise concurrency) or starved (producer/relays bound)?
val sat =
if (config.diagnose) {
" · ${activeWorkers.load()}/${config.drainConcurrency}w · ${throttled.load()} rl"
} else {
""
}
log(
"[graperank] round $progRound · ${human(roundDone.toLong())}/${human(progTarget.toLong())} ($pct%)" +
" · $rate/s · ~$eta · ${human(events)} ev · $parked slow · ${deadRelays.size()} dead",
" · $rate/s · ~$eta · ${human(events)} ev · $parked slow · ${deadRelays.size()} dead$sat",
)
}
}
}
}
/**
* A drain unit's page came back at the [FULL_PAGE_THRESHOLD] — it may have been
* truncated by the relay's per-REQ cap. Continue the SAME query in the
* background with `until` cursors ([fetchAllPages], starting at the page's
* oldest event, inclusive) to drain whatever the cap hid, streaming the extra
* events to [lateHarvest] just like a parked slow relay. Tracked by
* [parkedInFlight] so the round waits for it; gated by [limiter] and dropped if
* we have no [bgScope]. The boundary second is re-fetched and its already-seen
* events are dropped by [persist]'s crawl-wide dedup, so nothing double-counts.
* A no-op unless the page hit the threshold, so only dense units pay for it.
*/
private fun paginateIfCapped(
relay: NormalizedRelayUrl,
groupFilters: List<Filter>,
page: List<Pair<NormalizedRelayUrl, Event>>,
) {
if (page.size < FULL_PAGE_THRESHOLD) return
val scope = bgScope ?: return
val oldest = page.minOf { it.second.createdAt }
val contFilters = groupFilters.map { it.copy(until = oldest) }
parkedInFlight.addAndFetch(1)
scope.launch {
try {
val more = ArrayList<Pair<NormalizedRelayUrl, Event>>()
limiter.withPermit(relay) {
client.fetchAllPages(relay, contFilters, config.parkTimeoutMs) { ev -> more.add(relay to ev) }
}
for (pair in persist(more)) lateHarvest.trySend(pair)
} finally {
parkedInFlight.addAndFetch(-1)
}
}
}
/**
* Subscribe each relay to its filters behind [limiter] and drain them. A relay
* that reaches a terminal (EOSE/CLOSED/cannot-connect) within the FAST
@@ -766,11 +1167,13 @@ class GrapeRankDataCrawler(
relay: NormalizedRelayUrl,
into: ConcurrentMap<NormalizedRelayUrl, DrainFailure>,
) {
classifyDrainFailure(reason)?.let { kind ->
into.merge(relay, kind) { a, b ->
if (a == DrainFailure.HARD || b == DrainFailure.HARD) DrainFailure.HARD else DrainFailure.TRANSIENT
}
}
classifyDrainFailure(reason)?.let { kind -> into[relay] = kind }
}
// A relay asking us to slow down — the external concurrency ceiling.
fun isRateLimit(reason: String): Boolean {
val m = reason.lowercase()
return "429" in m || "too many" in m || "rate" in m || "throttl" in m
}
fun logSlow(
@@ -801,6 +1204,12 @@ class GrapeRankDataCrawler(
// this (conflated, so bursts collapse to one) and resets the park
// window, so a relay actively streaming is never cut mid-flight.
val activity = Channel<Unit>(Channel.CONFLATED)
// Latency breakdown: elapsed-since-[mark] of the first and last
// event, so the EOSE-wait AFTER the relay's last event (pure
// waiting on a done-but-slow relay) is separable from fetch time.
val mark = TimeSource.Monotonic.markNow()
val firstEvt = AtomicLong(-1)
val lastEvt = AtomicLong(-1)
val listener =
object : SubscriptionListener {
override fun onEvent(
@@ -811,6 +1220,9 @@ class GrapeRankDataCrawler(
) {
unitEvents.trySend(relay to event)
activity.trySend(Unit)
val e = mark.elapsedNow().inWholeMilliseconds
firstEvt.compareAndSet(-1, e)
lastEvt.store(e)
}
override fun onEose(
@@ -837,21 +1249,43 @@ class GrapeRankDataCrawler(
}
}
client.subscribe(subId, mapOf(subRelay to groupFilters), listener)
val mark = TimeSource.Monotonic.markNow()
val reason = withTimeoutOrNull(config.timeoutMs) { done.await() }
if (reason != null) {
// Terminal within the fast window — resolve this round.
val elapsedMs = mark.elapsedNow().inWholeMilliseconds
if (reason != "eose") notAnswered.add(subRelay)
classify(reason, subRelay, failures)
if (isRateLimit(reason)) throttled.addAndFetch(1)
// Split the wall time: fetch (to first event) vs EOSE-wait
// (after the last event) — only for drains that got events.
val firstE = firstEvt.load()
if (firstE >= 0) {
firstEventSumMs.addAndFetch(firstE)
eoseWaitSumMs.addAndFetch((elapsedMs - lastEvt.load()).coerceAtLeast(0))
drainWallSumMs.addAndFetch(elapsedMs)
drainSamples.addAndFetch(1)
}
if (elapsedMs > SLOW_DRAIN_LOG_MS) logSlow(subRelay, reason, elapsedMs, groupFilters)
unitEvents.close()
client.unsubscribe(subId)
val drained = buildList { for (e in unitEvents) add(e) }
telemetry.record(subRelay, RelayTelemetry.outcomeOf(reason, parked = false), elapsedMs, authorsIn(groupFilters), drained.size)
persist(drained)
val persisted = persist(drained)
// A full page from a clean EOSE may be the relay's cap, not the
// whole answer — background-paginate the remainder into lateHarvest.
if (reason == "eose") paginateIfCapped(subRelay, groupFilters, drained)
// Alive if it EOSE'd or handed us anything; a connect-timeout that
// gave nothing (classifyDrainFailure leaves it retryable forever)
// earns a strike toward eviction instead.
if (reason == "eose" || drained.isNotEmpty()) {
clearTimeoutStrikes(subRelay)
} else if (isTimeoutReason(reason)) {
strikeUnproductiveTimeout(subRelay)
}
persisted
} else {
// Still streaming — hand off and let the round move on.
burnedFastWindow.addAndFetch(1)
notAnswered.add(subRelay)
timedOut.add(subRelay)
val scope = bgScope
@@ -874,19 +1308,42 @@ class GrapeRankDataCrawler(
val lateDrained = buildList { for (e in unitEvents) add(e) }
telemetry.record(subRelay, RelayTelemetry.outcomeOf(late, parked = true), lateMs, authorsIn(groupFilters), lateDrained.size)
for (pair in persist(lateDrained)) lateHarvest.trySend(pair)
// A full parked page from a clean EOSE may also be capped —
// paginate its remainder in the background, same as the fast path.
if (late == "eose") paginateIfCapped(subRelay, groupFilters, lateDrained)
// Same liveness rule as the fast path: a park that ended
// in a clean EOSE or delivered anything clears the relay;
// one that idle-cut ("timeout") with nothing strikes it.
if (late == "eose" || lateDrained.isNotEmpty()) {
clearTimeoutStrikes(subRelay)
} else if (late == "timeout" || isTimeoutReason(late)) {
strikeUnproductiveTimeout(subRelay)
}
} finally {
client.unsubscribe(subId)
parkedInFlight.addAndFetch(-1)
}
}
// Parked: events are persisted into lateHarvest by the
// coroutine above, so this round contributes nothing here.
emptyList()
} else {
// Parking disabled (no bgScope, or parkTimeoutMs <= timeoutMs):
// drain and persist whatever streamed during the fast window
// instead of dropping it, then return it so the round ingests
// it exactly like the fast path — the other two branches persist,
// this one must too or those events are lost and re-queried.
val toMs = mark.elapsedNow().inWholeMilliseconds
logSlow(subRelay, "timeout", toMs, groupFilters)
telemetry.record(subRelay, RelayTelemetry.Outcome.FAST_TIMEOUT, toMs, authorsIn(groupFilters), 0)
unitEvents.close()
client.unsubscribe(subId)
val drained = buildList { for (e in unitEvents) add(e) }
telemetry.record(subRelay, RelayTelemetry.Outcome.FAST_TIMEOUT, toMs, authorsIn(groupFilters), drained.size)
// Nothing delivered → same unproductive-timeout signal, strike it;
// anything delivered proves the host productive and clears it.
if (drained.isEmpty()) strikeUnproductiveTimeout(subRelay) else clearTimeoutStrikes(subRelay)
persist(drained)
}
emptyList()
}
}
}
@@ -962,6 +1419,12 @@ class GrapeRankDataCrawler(
// Runs on [scope], so scope.cancel() at crawl end stops it.
scope.launch { progressTicker() }
// Background reachability culler: cheaply TCP-probes the cold tail of
// learned relays and drops the unreachable ones into deadHosts before the
// WS path pays the full connectTimeout on them. Runs on [scope], stopped
// by scope.cancel() at crawl end.
config.reachabilityProbe?.let { probe -> scope.launch { cullUnreachable(probe) } }
while (rounds < config.maxRounds) {
// Fold in whatever the parked (slow-but-alive) relays have delivered
// since the last round — their late contact lists expand the frontier
@@ -1006,6 +1469,10 @@ class GrapeRankDataCrawler(
// the majority cheaply (early rounds no-op until a backbone is learned).
val fedBeforeA = contactListsFed
shardedSweep(pending)
// Fold any kind:3 already sitting in the store — a prior round's
// ensureRelayLists co-fetch, a late parked delivery, or a previous run —
// so Phase B doesn't re-drain contact lists we already hold.
harvestFromStore(pending)
val phaseAMs = roundMark.elapsedNow().inWholeMilliseconds
val phaseAFed = contactListsFed - fedBeforeA
@@ -1017,7 +1484,7 @@ class GrapeRankDataCrawler(
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 }
val allLive = liveRelays.snapshot().filterTo(HashSet()) { !isDead(it) }
ensureRelayLists(stragglers.toSet(), allLive, scope)
// Continuous worker pool instead of chunked awaitAll barriers, so
@@ -1047,11 +1514,16 @@ class GrapeRankDataCrawler(
List(config.drainConcurrency) {
launch {
for ((batch, filters) in routed) {
val dead = HashMap<NormalizedRelayUrl, DrainFailure>()
val answered = HashSet<NormalizedRelayUrl>()
val events = drainGated(filters, dead, answered)
recordDead(dead)
drainedOut.send(DrainedBatch(batch, filters, answered, events))
activeWorkers.addAndFetch(1)
try {
val dead = HashMap<NormalizedRelayUrl, DrainFailure>()
val answered = HashSet<NormalizedRelayUrl>()
val events = drainGated(filters, dead, answered)
recordDead(dead)
drainedOut.send(DrainedBatch(batch, filters, answered, events))
} finally {
activeWorkers.addAndFetch(-1)
}
}
}
}
@@ -1110,6 +1582,12 @@ class GrapeRankDataCrawler(
// Crawl done — drop the warm pool.
client.unsubscribe(WARM_SUB_ID)
// Patient final pass: recover the stragglers the outbox model couldn't
// resolve by asking the content aggregators for their kind:3 ALONE, no
// longer racing the full fan-out (which cut the aggregators short during
// the rounds). Runs before [scope] is cancelled so slow aggregators park.
recoverStragglersFromAggregators()
// 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
@@ -1144,13 +1622,29 @@ class GrapeRankDataCrawler(
val stored = eventsStored.load()
log(
"[graperank] crawl complete: ${hopOf.size} discovered, $contactListsFed contact lists fed, " +
"${relaysContacted.size} relays contacted, ${deadRelays.size()} dead, $rounds rounds in $downloadMs ms; " +
"${relaysContacted.size} relays contacted, ${deadRelays.size()} dead + ${deadHosts.size()} timeout-evicted hosts, " +
"$rounds rounds in $downloadMs ms; " +
"by hop: " + hopHistogram.entries.joinToString(" ") { "${it.key}=${it.value}" },
)
log(
"[graperank] write path: $stored events stored, verify ${verifyMs}ms + insert ${insertMs}ms " +
"(summed across all drains, batch=${config.insertBatchSize})",
)
if (config.diagnose) {
val n = drainSamples.load().coerceAtLeast(1)
val wall = drainWallSumMs.load().coerceAtLeast(1)
val meanFirst = firstEventSumMs.load() / n
val meanEose = eoseWaitSumMs.load() / n
val eosePct = 100 * eoseWaitSumMs.load() / wall
log(
"[graperank] latency breakdown (drains with events, n=${drainSamples.load()}): " +
"mean time-to-first-event ${meanFirst}ms, mean EOSE-wait-after-last-event ${meanEose}ms " +
"($eosePct% of drain wall spent waiting for EOSE after the relay's last event); " +
"${burnedFastWindow.load()} drains blew the ${config.timeoutMs}ms fast window and parked; " +
"${throttled.load()} rate-limit responses. " +
"High EOSE-wait % → a shorter/adaptive fast window is the lever, not more concurrency.",
)
}
return Stats(
rounds = rounds,
contactListsFed = contactListsFed,
@@ -1161,6 +1655,10 @@ class GrapeRankDataCrawler(
verifyMs = verifyMs,
insertMs = insertMs,
eventsStored = stored,
// Only relays we actually dialed this run — exclude the seeded
// known-dead (skipped, not re-probed) so their original TTL stands.
deadRelays = deadRelays.snapshot() - config.knownDeadRelays,
liveRelays = liveRelays.snapshot(),
)
}
}
@@ -1356,6 +1854,24 @@ class GrapeRankDataCrawler(
// Authors per REQ filter — keeps individual subscriptions within relay limits.
private const val AUTHORS_PER_FILTER = 300
// Concurrent TCP reachability probes in the background culler. Raw sockets are
// cheap and short-lived; the per-relay WS limiter is unaffected (this never
// opens a REQ), so this only bounds file descriptors during the cull.
private const val PROBE_CONCURRENCY = 128
// Re-scan interval for the culler when it has probed everything learned so far
// and is waiting for new relays to be discovered.
private const val PROBE_IDLE_MS = 2000L
// A single REQ can match up to authors×kinds events; a relay that caps its
// response below that silently drops the tail (measured: user.kindpag.es
// returns at most ~100 events per REQ and ignores our limit). Any page that
// comes back with at least this many events is treated as possibly-capped and
// paginated with `until` cursors to drain the rest. Set at the smallest page
// cap we've observed, so it catches every relay that caps at or above it while
// sparing the common under-cap page an extra REQ.
private const val FULL_PAGE_THRESHOLD = 100
// 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.
@@ -1365,6 +1881,40 @@ class GrapeRankDataCrawler(
// timeout) is logged with its relay + filter, so slow relays can be replayed.
private const val SLOW_DRAIN_LOG_MS = 4000L
/**
* Is a drain terminal reason a connect/read TIMEOUT — the class
* [classifyDrainFailure] leaves retryable forever? The reason shape is
* `cannot:<message>` and the message now carries the exception class name
* (see BasicRelayClient), so a SocketTimeoutException surfaces as "timed
* out"/"timeout". Used to drive unproductive-timeout eviction.
*/
private fun isTimeoutReason(reason: String): Boolean {
if (!reason.startsWith("cannot")) return false
val m = reason.removePrefix("cannot:").lowercase()
return "timeout" in m || "timed out" in m
}
/**
* The authority (host[:port]) of a normalized relay URL — the segment between
* the `wss://` / `ws://` scheme and the first `/`. This is the key the
* timeout-eviction counts on, so the many per-user path URLs the outbox model
* mints for one server (`filter.nostr.wine/npubA`, `filter.nostr.wine/npubB`, …)
* collapse to a single evictable host. A bare host is its own authority, so this
* is a no-op for the common no-path relay. Deliberately host-only: it must NOT
* fold `filter.nostr.wine` into `nostr.wine` — those are different servers with
* different behaviour (the bare host may read fine while the filter host stalls).
*/
fun authorityOf(url: String): String {
val afterScheme =
when {
url.startsWith("wss://") -> url.substring(6)
url.startsWith("ws://") -> url.substring(5)
else -> url
}
val slash = afterScheme.indexOf('/')
return if (slash >= 0) afterScheme.substring(0, slash) else afterScheme
}
// Once the frontier is empty but parked relays are still streaming, how long
// to block waiting for one of them to deliver before re-checking convergence.
private const val PARK_POLL_MS = 2000L
@@ -1408,10 +1958,6 @@ class GrapeRankDataCrawler(
// 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
@@ -47,7 +47,7 @@ import kotlinx.coroutines.coroutineScope
* (it fell below the caller's cutoff, or dropped out of the graph) with a NIP-09
* kind:5 deletion, batched so the frame stays under the ~64KB event cap.
*
* Transport-agnostic like [GrapeRankDataCrawler]: it reads prior cards from an
* Transport-agnostic like [GrapeRankCrawler]: it reads prior cards from an
* [IEventStore] and emits through an injected [publish] function (event + relays →
* per-relay ack), so the store/relay wiring stays in the application while the
* reconcile + card-construction logic is reusable (e.g. by the Android app).
@@ -38,7 +38,7 @@ import com.vitorpamplona.quartz.nip65RelayList.AdvertisedRelayListEvent
* (kind:10002), and reports (kind:1984) of every author already known to the
* local [store].
*
* Where [GrapeRankDataCrawler] discovers the graph by walking follows outward from
* Where [GrapeRankCrawler] discovers the graph by walking follows outward from
* an observer, this refreshes what is *already* known: it reads every kind:10002 in
* the store, inverts them into a `write-relay -> authors` map (the outbox model — an
* author's events live on the relays they write to), fans those into one filter per
@@ -97,6 +97,14 @@ class GrapeRankUpdater(
val minAuthors: Int = 1,
val idleTimeoutMs: Long = 30_000L,
val publishTimeoutSecs: Long = 15,
/**
* Relays proven unreachable within the reachability cache's TTL (kind:30166).
* Skipped from the reconcile plan so we don't burn a connect timeout per dead
* relay — the crawl already found them dead, and a dead relay cannot serve its
* authors anyway. TTL'd, so a recovered relay is retried once the record ages
* out; this never drops a *live* author-advertised relay. See RelayReachabilityStore.
*/
val knownDead: Set<NormalizedRelayUrl> = emptySet(),
) {
/** Project the shared engine knobs onto a [NegentropyStoreSync.Config]. */
internal fun toEngineConfig() =
@@ -170,6 +178,7 @@ class GrapeRankUpdater(
val plan =
groups.entries
.filter { it.value.size >= config.minAuthors }
.filterNot { it.key in config.knownDead }
.sortedByDescending { it.value.size }
.map { it.key to it.value }
@@ -21,57 +21,48 @@
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).
* A drain per-relay failure worth acting on: the relay will not serve us THIS run,
* so drop it from further routing on the first occurrence. There is only one such
* verdict — [DEAD] — because re-probing hop-8's failed relays fresh, outside the
* crawl, showed the old "might clear, retry a few times" (TRANSIENT) bucket almost
* never clears: 503 Service Unavailable was 0/12 reachable, 502 Bad Gateway 3/15,
* connection-establishment failures 0/30, and the codes that WERE alive (403/402)
* are gated and will never hand us events. Spending extra dials on them was waste.
*
* - [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).
* The only two connect failures that genuinely recover are kept OUT of this verdict
* by [classifyDrainFailure] returning null (retry, never dead):
* - a **read** timeout — the relay accepted the handshake but is slow to serve;
* 12/18 (67%) were reachable fresh, only overloaded by the crawl's fan-out. The
* crawler's per-authority timeout strikes, which CLEAR on success, shed the gone.
* - an HTTP **429 / too many requests** — alive and rate-limiting; 4/4 reachable
* fresh. Retrying (spaced by the rate limiter) is how we eventually get its data.
*/
enum class DrainFailure { HARD, TRANSIENT }
enum class DrainFailure { DEAD, }
/**
* 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`.
* Classify a drain per-relay terminal reason. Returns null when the relay should be
* retried rather than dropped — a read/generic timeout, an alive 429 rate-limit, or
* a non-failure like eose/closed. Any other `cannot:<message>` (see
* `BasicRelayClient.onCannotConnect`) is [DrainFailure.DEAD]: it will not serve us
* this run, so drop it now instead of paying repeated connect attempts.
*/
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
// Alive, only asking us to slow down: an HTTP 429 / "too many requests" reliably
// clears — 4/4 such relays were reachable when re-probed fresh. Retry it (the
// rate limiter spaces our opens); never drop it.
if ("429" in m || "too many requests" in m) return null
// A READ timeout means the relay accepted the handshake but is slow to serve —
// 12/18 (67%) reachable fresh, alive but overloaded by the fan-out. Retry; the
// crawler's per-authority timeout strikes, which clear on success, shed the gone.
// A *connect* timeout is the opposite (the socket never opened, 0/30 reachable),
// so it is excluded here and falls through to DEAD with every other failure.
if (("timeout" in m || "timed out" in m) && "connect timed out" !in m) return null
// Everything else won't serve us this run: connect refused / unroutable / the
// proxy couldn't tunnel the CONNECT, a DNS or TLS failure, a dead-or-not-a-relay
// HTTP upgrade (502/503/500/504/410/404/200/…), or a mid-stream reset. Measured
// mostly dead (503 0%, 502 20% reachable) and, when alive, gated (402/403) or not
// a relay (200). Drop it now rather than burn more dials on it.
return DrainFailure.DEAD
}
@@ -0,0 +1,171 @@
/*
* 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.nip66RelayMonitor.reachability
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import com.vitorpamplona.quartz.nip01Core.signers.NostrSigner
import com.vitorpamplona.quartz.nip01Core.store.IEventStore
import com.vitorpamplona.quartz.nip66RelayMonitor.discovery.RelayDiscoveryEvent
import com.vitorpamplona.quartz.nip66RelayMonitor.discovery.networkType
import com.vitorpamplona.quartz.nip66RelayMonitor.discovery.rtt
import com.vitorpamplona.quartz.nip66RelayMonitor.discovery.tags.NetworkType
import com.vitorpamplona.quartz.nip66RelayMonitor.discovery.tags.RttType
import com.vitorpamplona.quartz.utils.TimeUtils
/**
* A durable, shareable relay-reachability cache backed by an [IEventStore] as
* NIP-66 **kind:30166 Relay Discovery** events — so the crawler, the WoT updater,
* and future runs all read and write the *same* liveness knowledge instead of each
* rediscovering dead relays from an in-memory set that is wiped when the process ends.
*
* ## Why NIP-66 / the event store
* A 30166 event is addressable by its `d`-tag (the normalized relay URL), so the
* store keeps exactly **one replaceable record per (monitor, relay)** — a natural
* per-relay status slot with a `created_at` timestamp that gives us a free TTL. The
* event store gives us persistence, cross-procedure sharing, and interop for free:
* 30166 events published by *other* monitors (nostr.watch et al.) can be ingested to
* seed reachability without probing, and our own records can be published back.
*
* ## How "dead" is represented
* NIP-66 has no explicit offline field; liveness is inferred from a fresh record that
* carries an `rtt-open` (a successful connection). This cache follows that convention:
* - **reachable** → a 30166 **with** `rtt-open`, `created_at` = probe time.
* - **dead** → a 30166 **without** `rtt-open` ("we checked, could not open"),
* `created_at` = probe time.
*
* So a fresh rtt-less record distinguishes *checked-and-dead* from *never-checked*
* (no record). When both a dead and a live record exist within the TTL for the same
* relay, **live wins** — any recent successful open overrides an earlier failure,
* whether the two came from us across time or from two different monitors.
*
* ## Not a replacement for the hot path
* [snapshot] is meant to be loaded ONCE at the start of a run into whatever in-memory
* structure the caller already uses for per-request `isDead` checks; [record] flushes
* a run's findings back at the end. It is deliberately not queried per routing decision.
*
* A relay is only ever skipped for the TTL window, never permanently — consistent with
* the outbox rule that every advertised write relay must be tried: a TTL'd record is
* "skip for now", not "ignore this author's home forever".
*
* ## The signer is a dedicated monitor service identity
* [signer] should be a **machine-level monitor key**, NOT a user/observer account: per
* NIP-66 a monitor is its own pubkey (which also publishes a kind:10166 announcement,
* a kind:0 profile and a kind:10002). Publishing these under the observer's key would
* conflate the WoT identity with a relay-monitoring service. [snapshot] still honours
* records from ANY author (so third-party monitors can be ingested); only [record]
* writes under this monitor key.
*/
class RelayReachabilityStore(
private val store: IEventStore,
private val signer: NostrSigner,
private val ttlSeconds: Long = DEFAULT_TTL_SECONDS,
) {
/**
* An in-memory view of the fresh (within-TTL) reachability records. [dead] holds
* relays proven unreachable and not since seen live; [live] holds relays with a
* recent successful open. A relay absent from both is simply unknown — re-probe it.
*/
class Snapshot(
val dead: Set<NormalizedRelayUrl>,
val live: Set<NormalizedRelayUrl>,
) {
fun isKnownDead(relay: NormalizedRelayUrl) = relay in dead
val size: Int get() = dead.size + live.size
}
/**
* Load every 30166 record fresher than [ttlSeconds] and fold it into a [Snapshot].
* Records from any monitor are honoured (live-wins), so ingesting third-party
* monitors' 30166 into [store] transparently improves the result.
*/
suspend fun snapshot(now: Long = TimeUtils.now()): Snapshot {
val since = now - ttlSeconds
val events =
store.query<RelayDiscoveryEvent>(
Filter(kinds = listOf(RelayDiscoveryEvent.KIND), since = since),
)
val live = HashSet<NormalizedRelayUrl>()
val dead = HashSet<NormalizedRelayUrl>()
for (ev in events) {
val relay = ev.relay() ?: continue
if (ev.rttOpen() != null) live.add(relay) else dead.add(relay)
}
// A recent successful open (from us later, or from another monitor) overrides
// an earlier dead mark for the same relay.
dead.removeAll(live)
return Snapshot(dead, live)
}
/**
* Persist a run's reachability findings as 30166 events: each [reachable] relay as
* a record WITH `rtt-open`, each [dead] relay (that is not also reachable) as one
* WITHOUT. Signed by [signer] and inserted into [store]; being addressable, each
* replaces this monitor's prior record for that relay, so the store stays bounded
* at roughly the number of distinct relays.
*
* [rttOpenMs] is the measured open round-trip in ms. It defaults to 0 as a **liveness
* flag only** — presence of the `rtt-open` tag, not its magnitude, is what [snapshot]
* reads as "reachable", and a caller that merely proved a relay served events (like
* the crawler) has no dedicated probe latency to report. A `0` therefore means
* "reachable, latency not probed by this writer", NOT a real 0 ms measurement. Do NOT
* publish these records to the wider network as authoritative latency data until a
* dedicated monitor probe supplies a real [rttOpenMs]; aggregators rank by it.
*/
suspend fun record(
reachable: Set<NormalizedRelayUrl>,
dead: Set<NormalizedRelayUrl>,
now: Long = TimeUtils.now(),
rttOpenMs: Long = 0,
) {
for (relay in reachable) writeOne(relay, up = true, now, rttOpenMs)
for (relay in dead) if (relay !in reachable) writeOne(relay, up = false, now, rttOpenMs)
}
private suspend fun writeOne(
relay: NormalizedRelayUrl,
up: Boolean,
now: Long,
rttOpenMs: Long,
) {
val template =
RelayDiscoveryEvent.build(relay, createdAt = now) {
networkType(networkTypeOf(relay))
if (up) rtt(RttType.OPEN, rttOpenMs)
}
store.insert(signer.sign(template))
}
companion object {
/** Default freshness window: a relay's status is trusted for a day, then re-probed. */
const val DEFAULT_TTL_SECONDS = 24L * 60 * 60
/** NIP-66 `n` network type inferred from the URL, so a `.onion`/i2p relay is tagged correctly. */
fun networkTypeOf(relay: NormalizedRelayUrl): NetworkType =
when {
RelayUrlNormalizer.isOnion(relay.url) -> NetworkType.TOR
relay.url.contains(".i2p") -> NetworkType.I2P
else -> NetworkType.CLEARNET
}
}
}
@@ -0,0 +1,66 @@
/*
* 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 kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNotEquals
/**
* [GrapeRankCrawler.authorityOf] is the key the crawl's timeout-eviction counts
* on. It must collapse the many per-user path URLs the outbox model mints for one
* server into a single host, WITHOUT folding a distinct sibling host (e.g. a
* `filter.` subdomain) into its parent.
*/
class GrapeRankAuthorityTest {
private fun auth(url: String) = GrapeRankCrawler.authorityOf(url)
@Test
fun bareHostIsItsOwnAuthority() {
assertEquals("relay.damus.io", auth("wss://relay.damus.io"))
assertEquals("relay.damus.io", auth("wss://relay.damus.io/"))
assertEquals("nos.lol", auth("ws://nos.lol"))
}
@Test
fun perUserPathUrlsOnOneHostCollapseToOneAuthority() {
val a = auth("wss://filter.nostr.wine/npub1aaaa?broadcast=true")
val b = auth("wss://filter.nostr.wine/npub1bbbb?broadcast=true&global=all")
val c = auth("wss://filter.nostr.wine/?global=all")
assertEquals("filter.nostr.wine", a)
assertEquals(a, b)
assertEquals(a, c)
}
@Test
fun filterSubdomainIsNotFoldedIntoBareHost() {
// nostr.wine reads are open; filter.nostr.wine is a different server that may
// stall — evicting one must never take out the other.
assertNotEquals(auth("wss://filter.nostr.wine/npub1x"), auth("wss://nostr.wine"))
}
@Test
fun portIsPartOfTheAuthority() {
assertEquals("relay.veganostr.com:443", auth("wss://relay.veganostr.com:443/npub1z"))
assertEquals("81.68.170.122:7114", auth("ws://81.68.170.122:7114/"))
assertNotEquals(auth("wss://example.com:443"), auth("wss://example.com:8080"))
}
}
@@ -0,0 +1,86 @@
/*
* Copyright (c) 2025 Vitor Pamplona
*
* Permission is hereby granted, free of charge, to any person obtaining a copy of
* this software and associated documentation files (the "Software"), to deal in
* the Software without restriction, including without limitation the rights to use,
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
* Software, and to permit persons to whom the Software is furnished to do so,
* subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all
* copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
*/
package com.vitorpamplona.quartz.nip01Core.relay.client.accessories
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNull
class DrainFailureTest {
// Non-failure and non-"cannot" terminals are never dead signals.
@Test
fun nonFailureTerminalsAreNull() {
assertNull(classifyDrainFailure("eose"))
assertNull(classifyDrainFailure("closed:duplicate: sub"))
assertNull(classifyDrainFailure("timeout"))
}
// A READ timeout (or generic post-handshake timeout) is alive-but-slow: never
// dead. Measured 67% of these relays were reachable when re-probed fresh.
@Test
fun readTimeoutsStayRetryable() {
assertNull(classifyDrainFailure("cannot:Read timed out (SocketTimeoutException)"))
assertNull(classifyDrainFailure("cannot:timeout (SocketTimeoutException)"))
}
// An HTTP 429 rate-limit is alive and will serve us after backoff: never dead.
// Measured 4/4 such relays reachable when re-probed fresh.
@Test
fun rateLimitStaysRetryable() {
assertNull(classifyDrainFailure("cannot:Server Misconfigured. Response: 429 Too Many Requests (ProtocolException)"))
}
// Failing to ESTABLISH the connection is dead (0/30 reachable fresh). "connect
// timed out" must be caught as DEAD and NOT slip into the read-timeout branch.
@Test
fun connectEstablishmentFailuresAreDead() {
assertEquals(DrainFailure.DEAD, classifyDrainFailure("cannot:Connect timed out (SocketTimeoutException)"))
assertEquals(DrainFailure.DEAD, classifyDrainFailure("cannot:Unexpected response code for CONNECT: (IOException)"))
assertEquals(DrainFailure.DEAD, classifyDrainFailure("cannot:Connection refused (ConnectException)"))
assertEquals(DrainFailure.DEAD, classifyDrainFailure("cannot:Failed to connect to /1.2.3.4:443"))
assertEquals(DrainFailure.DEAD, classifyDrainFailure("cannot:No route to host (NoRouteToHostException)"))
}
// DNS and TLS misconfig can never work: DEAD.
@Test
fun dnsAndTlsAreDead() {
assertEquals(DrainFailure.DEAD, classifyDrainFailure("cannot:Unable to resolve host (UnknownHostException)"))
assertEquals(DrainFailure.DEAD, classifyDrainFailure("cannot:Received fatal alert: unrecognized_name (SSLHandshakeException)"))
assertEquals(DrainFailure.DEAD, classifyDrainFailure("cannot:PKIX path building failed: certificate (CertificateException)"))
}
// Every other bad HTTP upgrade won't serve us this run (measured 503 0%, 502 20%
// reachable; 402/403 gated; 200 not a relay) — DEAD, dropped on the first strike.
@Test
fun deadOrGatedHttpUpgradesAreDead() {
assertEquals(DrainFailure.DEAD, classifyDrainFailure("cannot:Server Misconfigured. not a websocket"))
assertEquals(DrainFailure.DEAD, classifyDrainFailure("cannot:Server Misconfigured. Response: 503 Service Unavailable (ProtocolException)"))
assertEquals(DrainFailure.DEAD, classifyDrainFailure("cannot:Server Misconfigured. Response: 502 Bad Gateway (ProtocolException)"))
assertEquals(DrainFailure.DEAD, classifyDrainFailure("cannot:Server Misconfigured. Response: 402 Payment Required (ProtocolException)"))
}
// A mid-stream reset won't hand us events this run either: DEAD.
@Test
fun midStreamResetIsDead() {
assertEquals(DrainFailure.DEAD, classifyDrainFailure("cannot:Connection reset (SocketException)"))
assertEquals(DrainFailure.DEAD, classifyDrainFailure("cannot:Broken pipe (SocketException)"))
}
}
@@ -0,0 +1,117 @@
/*
* 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.nip66RelayMonitor.reachability
import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal
import com.vitorpamplona.quartz.nip01Core.store.sqlite.DefaultIndexingStrategy
import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
import com.vitorpamplona.quartz.nip66RelayMonitor.discovery.tags.NetworkType
import com.vitorpamplona.quartz.utils.Secp256k1Instance
import kotlinx.coroutines.runBlocking
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertTrue
class RelayReachabilityStoreTest {
private fun store() =
EventStore(
dbName = null,
indexStrategy = DefaultIndexingStrategy(),
)
private fun cache(store: EventStore) =
RelayReachabilityStore(
store = store,
signer = NostrSignerInternal(KeyPair()),
ttlSeconds = 3600,
)
private val live1 = RelayUrlNormalizer.normalize("wss://alive.example.com")
private val live2 = RelayUrlNormalizer.normalize("wss://also-alive.example.com")
private val dead1 = RelayUrlNormalizer.normalize("wss://dead.example.com")
private val dead2 = RelayUrlNormalizer.normalize("wss://gone.example.com")
private val onion = RelayUrlNormalizer.normalize("wss://abc.onion")
private val onionPath = RelayUrlNormalizer.normalize("wss://abc.onion/npub1x")
// Contains the literal ".onion" as a substring but is NOT a Tor host — a loose
// `contains(".onion")` would misclassify it; the normalizer's isOnion must not.
private val fakeOnion = RelayUrlNormalizer.normalize("wss://relay.onionfake.com")
@Test
fun recordsAndReloadsReachability() =
runBlocking {
Secp256k1Instance
val store = store()
val cache = cache(store)
val now = 1_000_000L
cache.record(reachable = setOf(live1, live2), dead = setOf(dead1, dead2), now = now)
val snap = cache.snapshot(now = now)
assertEquals(setOf(live1, live2), snap.live)
assertEquals(setOf(dead1, dead2), snap.dead)
assertTrue(snap.isKnownDead(dead1))
assertFalse(snap.isKnownDead(live1))
}
@Test
fun aFreshSuccessfulOpenOverridesAnEarlierDeadMark() =
runBlocking {
Secp256k1Instance
val store = store()
val cache = cache(store)
// Marked dead first, then seen alive a second later (addressable replace).
cache.record(reachable = emptySet(), dead = setOf(dead1), now = 1_000L)
cache.record(reachable = setOf(dead1), dead = emptySet(), now = 1_001L)
val snap = cache.snapshot(now = 1_001L)
assertTrue(dead1 in snap.live)
assertFalse(snap.isKnownDead(dead1))
}
@Test
fun recordsOlderThanTheTtlAreIgnored() =
runBlocking {
Secp256k1Instance
val store = store()
val cache = cache(store) // ttl = 3600s
cache.record(reachable = emptySet(), dead = setOf(dead1), now = 1_000L)
// "now" is well past the 1h TTL from when dead1 was recorded.
val snap = cache.snapshot(now = 1_000L + 3601L)
assertFalse(snap.isKnownDead(dead1))
assertEquals(0, snap.size)
}
@Test
fun onionRelayIsTaggedTorNetwork() {
assertEquals(NetworkType.TOR, RelayReachabilityStore.networkTypeOf(onion))
assertEquals(NetworkType.TOR, RelayReachabilityStore.networkTypeOf(onionPath))
assertEquals(NetworkType.CLEARNET, RelayReachabilityStore.networkTypeOf(live1))
// A host that merely contains ".onion" as a substring is clearnet, not Tor.
assertEquals(NetworkType.CLEARNET, RelayReachabilityStore.networkTypeOf(fakeOnion))
}
}