From a81561424724b7146345e39b397f5dbc1beeb790 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 14 Jul 2026 20:42:17 +0000 Subject: [PATCH] refactor(cli): move the relay census to `amy relay probe` MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The census probes the whole known relay universe and feeds the shared NIP-66 reachability cache (kind:30166) that every reachability-aware command reads — it was never graperank-specific, so it moves from GrapeRankCommand to RelayCommands as `amy relay probe`. `amy graperank probe` stays as an alias, and flags, JSON output, and behaviour are unchanged. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_013WSzVX9RoUxyV3nT3dfc56 --- cli/README.md | 1 + .../com/vitorpamplona/amethyst/cli/Main.kt | 11 ++- .../amethyst/cli/commands/GrapeRankCommand.kt | 89 ++----------------- .../amethyst/cli/commands/RelayCommands.kt | 88 +++++++++++++++++- 4 files changed, 101 insertions(+), 88 deletions(-) diff --git a/cli/README.md b/cli/README.md index bf393da508..91446e2e74 100644 --- a/cli/README.md +++ b/cli/README.md @@ -553,6 +553,7 @@ the last facet removes R entirely. | `amy relay add URL` / `remove URL` | Fan-out to the transport lists (nip65 `both` + `dm` + `key-package`). | | `amy relay list` | Print every configured relay bucket. | | `amy relay publish-lists` | Broadcast every configured relay list to the union of your relays. | +| `amy relay probe [--timeout SECS] [--concurrency N]` | The relay census: mass-connect every relay the local store knows (all stored kind:10002 relays + the reachability cache) in parallel waves and record live/dead + measured `rtt-open` into the NIP-66 reachability cache (kind:30166). Reachability-aware commands (`graperank crawl`/`update`) read it to skip dead relays and pre-connect live ones. (`amy graperank probe` remains as an alias.) | ### Local store maintenance diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Main.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Main.kt index d7410753bc..f4b2d7a5ef 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Main.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Main.kt @@ -484,6 +484,11 @@ private fun printUsage() { | relay list print every configured relay bucket | relay publish-lists broadcast every configured relay list | relay info URL fetch + print a relay's NIP-11 info document + | relay probe [--timeout SECS] relay census: mass-connect every relay the store + | [--concurrency N] knows and record live/dead + measured rtt-open + | into the reachability cache (NIP-66 kind:30166), + | so reachability-aware commands (graperank crawl/ + | update) skip dead relays and wait once | outbox USER [--refresh] show USER's NIP-65 read/write relays (outbox model) | [--timeout SECS] (USER: npub|nprofile|hex|name@domain) | @@ -639,10 +644,8 @@ private fun printUsage() { | [--max-hops N] [--preconnect-cap N] 1984/10002) into the local store without scoring. | [--no-preconnect] Pre-connects every known-live relay in one parallel | storm (seeded from the reachability cache). - | graperank probe [--timeout SECS] relay census: mass-connect every relay the store - | [--concurrency N] knows and record live/dead + measured rtt-open into - | the reachability cache (NIP-66 kind:30166), so the - | next crawl skips dead relays and waits once. + | graperank probe alias for `relay probe` (the census moved there — + | it feeds the shared NIP-66 reachability cache). | graperank update [--down] [--up] refresh every locally-known author's WoT record kinds | [--no-sync-deletions] [--timeout SECS] (0/3/10002/1984) from their own outbox: reads all | [--relay-concurrency N] [--author-chunk N] kind:10002 in the store, groups authors by write 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 124756afae..ac836681fb 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 @@ -44,7 +44,6 @@ import com.vitorpamplona.quartz.nip09Deletions.DeletionEvent import com.vitorpamplona.quartz.nip09Deletions.DeletionIndex import com.vitorpamplona.quartz.nip51Lists.muteList.MuteListEvent import com.vitorpamplona.quartz.nip56Reports.ReportEvent -import com.vitorpamplona.quartz.nip66RelayMonitor.reachability.RelayProber import com.vitorpamplona.quartz.nip85TrustedAssertions.list.TrustProviderListEvent import com.vitorpamplona.quartz.nip85TrustedAssertions.list.serviceProviders import com.vitorpamplona.quartz.nip85TrustedAssertions.list.tags.ProviderTypes @@ -92,9 +91,10 @@ import kotlin.math.roundToInt * relay(s) converge to the local card set via a NIP-77 up-only reconcile * (nothing is re-signed or re-scored), and refresh the observer's kind:10040 * pointer when we hold their key. - * - `amy graperank probe` — the relay census: mass-connect every relay the store + * - `amy relay probe` — the relay census: mass-connect every relay the store * knows and record live/dead + measured RTT into the reachability cache, so the * next crawl skips the dead and pre-connects the living in one parallel storm. + * Lives in [RelayCommands] (`graperank probe` is kept as an alias). * - bare `amy graperank [OBSERVER]` — the convenience combo: crawl then score. * * Sub-verbs complete the NIP-85 experience — discovery and consumption: @@ -201,7 +201,10 @@ object GrapeRankCommand { // `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()) - "probe" -> probe(dataDir, tail.drop(1).toTypedArray()) + // The relay census outgrew graperank (it feeds the shared NIP-66 + // reachability cache every command reads) and moved to `amy relay + // probe`; this alias keeps the old spelling working. + "probe" -> RelayCommands.probe(dataDir, tail.drop(1).toTypedArray()) "update" -> update(dataDir, tail.drop(1).toTypedArray()) "score" -> run(dataDir, tail.drop(1).toTypedArray(), forceOffline = true) "publish" -> publish(dataDir, tail.drop(1).toTypedArray()) @@ -540,86 +543,6 @@ object GrapeRankCommand { return 0 } - /** - * `amy graperank probe [--timeout SECS] [--concurrency N]` — - * the relay census. Mass-connects the ENTIRE relay universe the local store knows - * (every relay advertised in any stored kind:10002, deduped per host, plus - * everything already in the reachability cache) in parallel waves with a no-op - * REQ, so the "is this relay alive, and how slow?" wait is paid once, up front, - * concurrently — then records per-relay verdicts with real measured `rtt-open` - * into the NIP-66 reachability cache (kind:30166). - * - * The next `graperank crawl` reads that cache to (a) skip the dead set without - * dialing it and (b) pre-connect the live set in one storm — separating "working - * but slow" (kept; the crawler's patient park path waits for them) from "not - * working" (skipped entirely). Typical flow the first time: - * `graperank crawl --max-hops 2` (cheap, saves the relay lists) → `graperank - * probe` → full `graperank crawl`. - */ - private suspend fun probe( - dataDir: DataDir, - rest: Array, - ): Int { - val args = Args(rest) - val timeoutMs = args.longFlag("timeout", 15L) * 1000 - val waveSize = args.intFlag("concurrency", Context.defaultPreconnectCap) - - Context.openOrAnonymous(dataDir).use { ctx -> - ctx.prepare() - val cached = ctx.reachability.snapshot() - val universe = RelayProber.knownRelayUniverse(ctx.store) + cached.live + cached.dead - if (universe.isEmpty()) { - Output.emit( - linkedMapOf( - "probed" to 0, - "note" to "no relays known locally — run `amy graperank crawl` first to gather kind:10002 relay lists", - ), - ) - return 0 - } - - System.err.println( - "[relay-probe] probing ${universe.size} relays in waves of $waveSize " + - "(${timeoutMs / 1000}s per wave; open-files limit ${Context.maxFileDescriptors})", - ) - val result = - RelayProber(ctx.client) { System.err.println(it) } - .probe(universe, timeoutMs, waveSize) - - ctx.reachability.recordProbed(result.reachableRttMs(), result.deadRelays()) - - val rtts = - result.reachable - .map { it.rttOpenMs } - .filter { it >= 0 } - .sorted() - - fun pct(p: Int): Long? = if (rtts.isEmpty()) null else rtts[(rtts.size - 1) * p / 100] - val slowest = - result.reachable - .filter { it.rttOpenMs >= 0 } - .sortedByDescending { it.rttOpenMs } - .take(10) - .map { mapOf("relay" to it.relay.url, "rtt_open_ms" to it.rttOpenMs) } - val authWalled = result.reachable.count { it.error?.startsWith("closed:") == true } - - Output.emit( - linkedMapOf( - "probed" to result.verdicts.size, - "reachable" to result.reachable.size, - "dead" to result.dead.size, - "closed_by_policy" to authWalled, - "elapsed_ms" to result.elapsedMs, - "rtt_open_p50_ms" to pct(50), - "rtt_open_p90_ms" to pct(90), - "rtt_open_p99_ms" to pct(99), - "slowest" to slowest, - ), - ) - } - return 0 - } - /** * `amy graperank update [flags]` — refresh every locally-known author's WoT * record kinds (0 / 3 / 10002 / 1984) straight from their own outbox, so the diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/RelayCommands.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/RelayCommands.kt index 74f48c1506..33207c6790 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/RelayCommands.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/RelayCommands.kt @@ -43,6 +43,7 @@ import com.vitorpamplona.quartz.nip51Lists.relayLists.TrustedRelayListEvent import com.vitorpamplona.quartz.nip65RelayList.AdvertisedRelayListEvent import com.vitorpamplona.quartz.nip65RelayList.tags.AdvertisedRelayInfo import com.vitorpamplona.quartz.nip65RelayList.tags.AdvertisedRelayType +import com.vitorpamplona.quartz.nip66RelayMonitor.reachability.RelayProber import okhttp3.OkHttpClient import okhttp3.Request @@ -86,7 +87,7 @@ import okhttp3.Request */ object RelayCommands { private const val USAGE = - "relay …" + "relay …" // ------------------------------------------------------------------ // Flat buckets — a plain list of relay URLs, one Nostr replaceable kind. @@ -216,6 +217,7 @@ object RelayCommands { // `info` is also intercepted in Main before account resolution // (it needs no account); routed here too for when one exists. "info" -> info(rest) + "probe" -> probe(dataDir, rest) "add" -> fanOut(dataDir, Args(rest), add = true) "remove", "rm" -> fanOut(dataDir, Args(rest), add = false) "outbox" -> facetVerb(dataDir, Facet.OUTBOX, rest) @@ -269,6 +271,90 @@ object RelayCommands { } } + // ------------------------------------------------------------------ + // relay probe — the relay census (feeds the NIP-66 reachability cache) + // ------------------------------------------------------------------ + + /** + * `amy relay probe [--timeout SECS] [--concurrency N]` — + * the relay census. Mass-connects the ENTIRE relay universe the local store knows + * (every relay advertised in any stored kind:10002, deduped per host, plus + * everything already in the reachability cache) in parallel waves with a no-op + * REQ, so the "is this relay alive, and how slow?" wait is paid once, up front, + * concurrently — then records per-relay verdicts with real measured `rtt-open` + * into the NIP-66 reachability cache (kind:30166). + * + * Every reachability-aware command reads that cache to skip the dead set without + * dialing it: `graperank crawl` also pre-connects the live set in one storm, + * separating "working but slow" (kept; the crawler's patient park path waits for + * them) from "not working" (skipped entirely). Typical flow the first time: + * `graperank crawl --max-hops 2` (cheap, saves the relay lists) → `relay probe` + * → full `graperank crawl`. (`graperank probe` is kept as an alias.) + */ + suspend fun probe( + dataDir: DataDir, + rest: Array, + ): Int { + val args = Args(rest) + val timeoutMs = args.longFlag("timeout", 15L) * 1000 + val waveSize = args.intFlag("concurrency", Context.defaultPreconnectCap) + + Context.openOrAnonymous(dataDir).use { ctx -> + ctx.prepare() + val cached = ctx.reachability.snapshot() + val universe = RelayProber.knownRelayUniverse(ctx.store) + cached.live + cached.dead + if (universe.isEmpty()) { + Output.emit( + linkedMapOf( + "probed" to 0, + "note" to "no relays known locally — run `amy graperank crawl` first to gather kind:10002 relay lists", + ), + ) + return 0 + } + + System.err.println( + "[relay-probe] probing ${universe.size} relays in waves of $waveSize " + + "(${timeoutMs / 1000}s per wave; open-files limit ${Context.maxFileDescriptors})", + ) + val result = + RelayProber(ctx.client) { System.err.println(it) } + .probe(universe, timeoutMs, waveSize) + + ctx.reachability.recordProbed(result.reachableRttMs(), result.deadRelays()) + + val rtts = + result.reachable + .map { it.rttOpenMs } + .filter { it >= 0 } + .sorted() + + fun pct(p: Int): Long? = if (rtts.isEmpty()) null else rtts[(rtts.size - 1) * p / 100] + val slowest = + result.reachable + .filter { it.rttOpenMs >= 0 } + .sortedByDescending { it.rttOpenMs } + .take(10) + .map { mapOf("relay" to it.relay.url, "rtt_open_ms" to it.rttOpenMs) } + val authWalled = result.reachable.count { it.error?.startsWith("closed:") == true } + + Output.emit( + linkedMapOf( + "probed" to result.verdicts.size, + "reachable" to result.reachable.size, + "dead" to result.dead.size, + "closed_by_policy" to authWalled, + "elapsed_ms" to result.elapsedMs, + "rtt_open_p50_ms" to pct(50), + "rtt_open_p90_ms" to pct(90), + "rtt_open_p99_ms" to pct(99), + "slowest" to slowest, + ), + ) + } + return 0 + } + // ------------------------------------------------------------------ // Flat-bucket verbs // ------------------------------------------------------------------