mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
refactor(cli): move the relay census to amy relay probe
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 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_013WSzVX9RoUxyV3nT3dfc56
This commit is contained in:
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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<String>,
|
||||
): 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<String, Any?>(
|
||||
"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<String, Any?>(
|
||||
"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
|
||||
|
||||
@@ -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 <outbox|inbox|nip65|dm|key-package|search|private|blocked|trusted|proxy|indexer|broadcast|feeds|add|remove|list|publish-lists|info> …"
|
||||
"relay <outbox|inbox|nip65|dm|key-package|search|private|blocked|trusted|proxy|indexer|broadcast|feeds|add|remove|list|publish-lists|info|probe> …"
|
||||
|
||||
// ------------------------------------------------------------------
|
||||
// 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<String>,
|
||||
): 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<String, Any?>(
|
||||
"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<String, Any?>(
|
||||
"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
|
||||
// ------------------------------------------------------------------
|
||||
|
||||
Reference in New Issue
Block a user