diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt index 4643c406ea..e74998b91d 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt @@ -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 = 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) } diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/OperatorKeys.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/OperatorKeys.kt index 71a7dd6d68..68c095b234 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/OperatorKeys.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/OperatorKeys.kt @@ -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:" } } diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/GrapeRankCommand.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/GrapeRankCommand.kt index c41c71b999..e9a7177e36 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/GrapeRankCommand.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/GrapeRankCommand.kt @@ -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 = + 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? = + 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, @@ -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, ): 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( @@ -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 diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankDataCrawler.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankCrawler.kt similarity index 65% rename from quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankDataCrawler.kt rename to quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankCrawler.kt index cbc47183d9..63ebda32bf 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankDataCrawler.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankCrawler.kt @@ -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, val contentFallbackRelays: Set, + val contentAggregatorRelays: Set = 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 = 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, + val liveRelays: Set, ) /** @@ -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() + + // 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() + + // 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() val relaysContacted = hashSetOf() - val writeRelayFreq = HashMap() - val liveRelays = hashSetOf() + val writeRelayFreq = ConcurrentMap() + val liveRelays = ConcurrentSet() // 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>() val attempts = ConcurrentMap() - val deadRelays = ConcurrentSet() - val relayStrikes = ConcurrentMap() + + // 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().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() + val deadHosts = ConcurrentSet() + + // 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() // 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) { - 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 = - 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() // 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): 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, bgScope: CoroutineScope, ) { - val missing = pubkeys.filter { relaysOf(it) == null } - if (missing.isEmpty()) return - suspend fun query( authors: List, relays: Set, + kinds: List, ) { 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>() 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>) { + 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, + page: List>, + ) { + 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>() + 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, ) { - 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(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() - val answered = HashSet() - val events = drainGated(filters, dead, answered) - recordDead(dead) - drainedOut.send(DrainedBatch(batch, filters, answered, events)) + activeWorkers.addAndFetch(1) + try { + val dead = HashMap() + val answered = HashSet() + 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:` 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 diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankPublisher.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankPublisher.kt index cd03bd7d04..1ccc26cea2 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankPublisher.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankPublisher.kt @@ -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). diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankUpdater.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankUpdater.kt index 5d6d1cc85a..e84778f4f5 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankUpdater.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankUpdater.kt @@ -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 = 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 } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/DrainFailure.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/DrainFailure.kt index b2d605d044..1786ac04c2 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/DrainFailure.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/DrainFailure.kt @@ -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:` 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:` (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 } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip66RelayMonitor/reachability/RelayReachabilityStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip66RelayMonitor/reachability/RelayReachabilityStore.kt new file mode 100644 index 0000000000..93f6648639 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip66RelayMonitor/reachability/RelayReachabilityStore.kt @@ -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, + val live: Set, + ) { + 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( + Filter(kinds = listOf(RelayDiscoveryEvent.KIND), since = since), + ) + val live = HashSet() + val dead = HashSet() + 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, + dead: Set, + 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 + } + } +} diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankAuthorityTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankAuthorityTest.kt new file mode 100644 index 0000000000..e43aaf0e02 --- /dev/null +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/experimental/graperank/GrapeRankAuthorityTest.kt @@ -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")) + } +} diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/DrainFailureTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/DrainFailureTest.kt new file mode 100644 index 0000000000..cf1c993344 --- /dev/null +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/DrainFailureTest.kt @@ -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)")) + } +} diff --git a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip66RelayMonitor/reachability/RelayReachabilityStoreTest.kt b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip66RelayMonitor/reachability/RelayReachabilityStoreTest.kt new file mode 100644 index 0000000000..44ce0df4fc --- /dev/null +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip66RelayMonitor/reachability/RelayReachabilityStoreTest.kt @@ -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)) + } +}