perf(cli): adaptive per-relay subscription cap for the crawl

Replace the blunt global concurrency number with per-relay back-pressure.
Every relay starts generous (100 concurrent subscriptions) and is demoted
down a ladder (100 -> 20 -> 10) only when it complains about concurrency —
a CLOSED rate-limited, or a NOTICE like "too many concurrent REQs" /
"too many subscriptions" / "burst exhausted". Well-behaved relays keep
the full cap; only the busy hubs that push back get throttled, and only as
far as they keep pushing.

AdaptiveRelayLimiter registers as a RelayConnectionListener so demotions
are driven straight off the same NOTICE/CLOSED frames RelayDiagnostics
already observes, keyed by relay.url. Context.drain gains a gatePerRelay
path that opens one subscription per relay, each held behind that relay's
gate, so our concurrent subs on it never exceed its current cap. The gate
is a fair FIFO bounded semaphore whose limit can only be lowered; shrinking
below the in-use count admits no new subs until enough finish, so
concurrency converges down to the new cap.

Because a hot relay can no longer be flooded, the global content-drain
fan-out is raised (18 -> 48) to crawl the many well-behaved relays faster.
The crawl emits a relay_throttling summary (which relays were capped, and
to what) alongside relay_feedback.
This commit is contained in:
Claude
2026-07-07 11:44:54 +00:00
parent fb0011129d
commit de8648f3f1
3 changed files with 332 additions and 14 deletions
@@ -0,0 +1,198 @@
/*
* 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.amethyst.cli
import com.vitorpamplona.quartz.nip01Core.relay.client.listeners.RelayConnectionListener
import com.vitorpamplona.quartz.nip01Core.relay.client.single.IRelayClient
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.ClosedMessage
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.NoticeMessage
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.atomic.AtomicInteger
/**
* Adaptive per-relay concurrent-subscription cap.
*
* Every relay starts with a generous cap ([startCap], default 100) — we assume a
* relay can take as many concurrent REQs as we throw at it until it tells us
* otherwise. When a relay complains about concurrency (a `CLOSED rate-limited`,
* or a `NOTICE` like "too many concurrent REQs" / "too many subscriptions" /
* "burst exhausted"), we demote *that relay only* down the [ladder]
* (100 → 20 → 10). A well-behaved relay keeps the full cap; only the ones that
* push back get throttled, and only as far as they keep pushing.
*
* This replaces a single blunt global concurrency number with per-relay
* back-pressure: the crawl can fan out widely across the many relays that don't
* mind, while automatically easing off the few busy hubs that do — exactly the
* signals [RelayDiagnostics] already observes, here turned into an actuator.
*
* Registered as a [RelayConnectionListener] on the shared client, so demotions
* are driven straight off the incoming NOTICE/CLOSED frames (which fire on the
* per-relay socket threads — all state here is concurrent). Drains gate through
* [withPermit]; [Context.drain]'s `gatePerRelay` path holds a relay's permit for
* the lifetime of that relay's subscription, so at most `cap` of our
* subscriptions are ever open on it at once.
*/
class AdaptiveRelayLimiter(
private val startCap: Int = 100,
private val ladder: List<Int> = listOf(20, 10),
) : RelayConnectionListener {
private val gates = ConcurrentHashMap<NormalizedRelayUrl, Gate>()
// How many concurrency complaints we've acted on per relay (== index+1 into
// the ladder). Capped at ladder.size: past the floor we stop demoting.
private val demotions = ConcurrentHashMap<NormalizedRelayUrl, Int>()
private fun gate(relay: NormalizedRelayUrl): Gate = gates.getOrPut(relay) { Gate(startCap) }
/** Run [block] holding one of [relay]'s permits, respecting its current cap. */
suspend fun <T> withPermit(
relay: NormalizedRelayUrl,
block: suspend () -> T,
): T {
val g = gate(relay)
g.acquire()
try {
return block()
} finally {
g.release()
}
}
override fun onIncomingMessage(
relay: IRelayClient,
msgStr: String,
msg: Message,
) {
when (msg) {
is ClosedMessage -> if (isConcurrencyComplaint(msg.message)) demote(relay.url)
is NoticeMessage -> if (isConcurrencyComplaint(msg.message)) demote(relay.url)
else -> Unit
}
}
/** Step [relay] one rung down the cap ladder, unless it's already at the floor. */
private fun demote(relay: NormalizedRelayUrl) {
// Fast path: relays flood identical NOTICEs, so bail once at the floor
// instead of counting them all (the demotion is monotonic and idempotent).
if ((demotions[relay] ?: 0) >= ladder.size) return
val step = demotions.merge(relay, 1, Int::plus)!!
val cap = ladder[(step - 1).coerceIn(0, ladder.size - 1)]
gate(relay).lower(cap)
if (step <= ladder.size) {
System.err.println("[limiter] ${relay.url} capped at $cap concurrent subs (complaint #$step)")
}
}
private fun isConcurrencyComplaint(text: String): Boolean {
val t = text.lowercase()
return CONCURRENCY_MARKERS.any { it in t }
}
/** JSON-friendly view of which relays we throttled and how far. */
fun snapshot(): Map<String, Any?> {
val cappedAt = sortedMapOf<Int, Int>()
for ((_, step) in demotions) {
val cap = ladder[(step - 1).coerceIn(0, ladder.size - 1)]
cappedAt.merge(cap, 1, Int::plus)
}
return mapOf(
"start_cap" to startCap,
"ladder" to ladder,
"throttled_relays" to demotions.size,
"capped_at" to cappedAt,
)
}
fun hadThrottling(): Boolean = demotions.isNotEmpty()
/**
* A bounded-concurrency gate whose limit can only ever be *lowered* (relays
* never earn their cap back within a run). Fair FIFO hand-off: a released
* permit goes to the longest-waiting acquirer. Lowering the limit below the
* in-use count doesn't cancel live holders — it just refuses to admit new
* ones until enough release that `inUse < limit` again, so the concurrency
* converges down to the new cap as the excess subscriptions finish.
*/
private class Gate(
initialLimit: Int,
) {
private val limit = AtomicInteger(initialLimit)
private val mutex = Mutex()
private var inUse = 0
private val waiters = ArrayDeque<CompletableDeferred<Unit>>()
suspend fun acquire() {
val wait =
mutex.withLock {
if (inUse < limit.get()) {
inUse++
null
} else {
CompletableDeferred<Unit>().also { waiters.addLast(it) }
}
}
wait?.await()
}
suspend fun release() {
mutex.withLock {
inUse--
while (inUse < limit.get() && waiters.isNotEmpty()) {
waiters.removeFirst().complete(Unit)
inUse++
}
}
}
/** Monotonically shrink the cap. Safe to call from any thread. */
fun lower(newLimit: Int) {
limit.updateAndGet { if (newLimit < it) newLimit else it }
}
}
companion object {
// Substrings (matched case-insensitively) that mean "you're opening too
// many concurrent subscriptions / sending too fast" — the failure modes a
// lower per-relay cap actually fixes. Auth/blocked/restricted/unsupported
// are deliberately excluded: throttling wouldn't help those.
private val CONCURRENCY_MARKERS =
listOf(
"too many concurrent",
"concurrent req",
"too many subscription",
"number of subscriptions",
"subscription limit",
"too many req",
"rate-limit",
"rate limit",
"ratelimit",
"burst exhausted",
"throttl",
"too many messages",
"slow down",
)
}
}
@@ -74,10 +74,12 @@ import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.channels.Channel.Factory.UNLIMITED
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.joinAll
import kotlinx.coroutines.launch
import kotlinx.coroutines.selects.select
import kotlinx.coroutines.withTimeoutOrNull
import okhttp3.OkHttpClient
import java.util.concurrent.ConcurrentHashMap
/**
* Per-invocation wiring. Each CLI run constructs a Context, does its work,
@@ -154,6 +156,16 @@ class Context(
*/
val relayDiagnostics: RelayDiagnostics = RelayDiagnostics().also { client.addConnectionListener(it) }
/**
* Adaptive per-relay concurrent-subscription cap. Starts every relay
* generous (100) and demotes only the ones that complain about concurrency
* (100 → 20 → 10), driven straight off the NOTICE/CLOSED frames it observes
* as a connection listener. [drain]'s `gatePerRelay` path holds a relay's
* permit for the life of that relay's subscription, so we never exceed the
* cap the relay itself asked for. Idle for commands that don't opt in.
*/
val relayLimiter: AdaptiveRelayLimiter = AdaptiveRelayLimiter().also { client.addConnectionListener(it) }
/**
* NIP-42 responder: answers a relay's AUTH challenge by signing with the
* account key, so auth-gated relays serve our reads instead of CLOSing the
@@ -441,8 +453,10 @@ class Context(
timeoutMs: Long = 8_000,
diagnoseSlow: Boolean = false,
deadOut: MutableSet<NormalizedRelayUrl>? = null,
gatePerRelay: Boolean = false,
): List<Pair<NormalizedRelayUrl, Event>> {
if (filters.isEmpty()) return emptyList()
if (gatePerRelay) return drainGated(filters, timeoutMs, diagnoseSlow, deadOut)
val eventChannel = Channel<Pair<NormalizedRelayUrl, Event>>(UNLIMITED)
// Carries the terminal reason per relay so a timeout can distinguish a slow
// relay (never terminal) from a connect failure / CLOSED.
@@ -525,6 +539,107 @@ class Context(
return collected
}
/**
* Per-relay-gated variant of [drain] used by the crawl. Instead of one
* subscription spanning every relay, each relay gets its own subscription
* held behind [relayLimiter], so we never exceed the relay's adaptive
* concurrent-subscription cap. A relay whose cap is full simply waits for one
* of our other subscriptions on it to finish before its REQ goes out; relays
* we haven't upset run at the full starting cap and never wait.
*
* Semantics match [drain] otherwise: verify+store on a single consumer
* (so store writes stay serialized), return events tagged by relay, and
* report hard connect failures into [deadOut].
*/
private suspend fun drainGated(
filters: Map<NormalizedRelayUrl, List<Filter>>,
timeoutMs: Long,
diagnoseSlow: Boolean,
deadOut: MutableSet<NormalizedRelayUrl>?,
): List<Pair<NormalizedRelayUrl, Event>> {
val eventChannel = Channel<Pair<NormalizedRelayUrl, Event>>(UNLIMITED)
// One relay per subId, so the relay alone identifies which subscription a
// callback is for. First terminal frame wins; a timeout leaves it unset.
val relayDone = ConcurrentHashMap<NormalizedRelayUrl, CompletableDeferred<String>>()
for (r in filters.keys) relayDone[r] = CompletableDeferred()
val doneReasons = ConcurrentHashMap<NormalizedRelayUrl, String>()
val listener =
object : SubscriptionListener {
override fun onEvent(
event: Event,
isLive: Boolean,
relay: NormalizedRelayUrl,
forFilters: List<Filter>?,
) {
eventChannel.trySend(relay to event)
}
override fun onEose(
relay: NormalizedRelayUrl,
forFilters: List<Filter>?,
) {
relayDone[relay]?.complete("eose")
}
override fun onClosed(
message: String,
relay: NormalizedRelayUrl,
forFilters: List<Filter>?,
) {
relayDone[relay]?.complete("closed:$message")
}
override fun onCannotConnect(
relay: NormalizedRelayUrl,
message: String,
forFilters: List<Filter>?,
) {
relayDone[relay]?.complete("cannot:$message")
}
}
val collected = mutableListOf<Pair<NormalizedRelayUrl, Event>>()
coroutineScope {
// Single consumer: verify+store serially, exactly like drain().
val consumer =
launch {
for ((relay, event) in eventChannel) {
if (verifyAndStore(event)) collected.add(relay to event)
}
}
// One gated subscription per relay. The permit is held for the whole
// life of the relay's REQ, so concurrent subs on it never exceed its cap.
filters
.map { (relay, relayFilters) ->
launch {
relayLimiter.withPermit(relay) {
val subId = newSubId()
client.subscribe(subId, mapOf(relay to relayFilters), listener)
try {
val reason = withTimeoutOrNull(timeoutMs) { relayDone[relay]!!.await() }
doneReasons[relay] = reason ?: "timeout"
} finally {
client.unsubscribe(subId)
}
}
}
}.joinAll()
// All subscriptions are torn down; no more events can arrive. Close the
// channel so the consumer drains what's buffered and completes.
eventChannel.close()
consumer.join()
}
if (diagnoseSlow) {
val stalled = filters.keys.filter { (doneReasons[it] ?: "timeout") == "timeout" }.toSet()
if (stalled.isNotEmpty()) logSlowDrain(timeoutMs, stalled, doneReasons, collected)
}
deadOut?.let { out ->
for ((relay, reason) in doneReasons) {
if (reason.startsWith("cannot")) out.add(relay)
}
}
return collected
}
/**
* On a [drain] timeout, report which relays stalled and why — a relay that
* never sent EOSE (slow, possibly still streaming) vs one that couldn't be
@@ -94,16 +94,17 @@ object GrapeRankCommand {
// (empirically ~250 users/drain succeeds, ~17k fails); keep the fan-out small.
private const val USER_BATCH = 256
// Concurrent content drains. Each drain opens exactly ONE subscription per
// relay it touches (one REQ per relay under the drain's subId), so this IS
// the per-relay concurrency cap: a popular relay shared by many pending users
// receives at most this many concurrent subs from us — while the fan-out
// across *different* relays stays fully parallel (each drain hits ~100 distinct
// outboxes). RelayDiagnostics showed 24 blew past typical relay limits
// (rate-limited=1433, "too many concurrent REQs"=1286, "too many
// subscriptions"=710). Target ~20 subs/relay; 18 leaves room for the 1
// persistent warm-pool sub so the peak stays ~19, just under the common cap.
private const val DRAIN_CONCURRENCY = 18
// Global content-drain fan-out — how many outbox batches we drain at once.
// This is now purely a GLOBAL bound (memory / open sockets); the per-relay
// concurrency limit is enforced separately and adaptively by
// [AdaptiveRelayLimiter] (drains run with gatePerRelay=true), which starts
// every relay at 100 concurrent subs and demotes only the ones that complain
// (100 → 20 → 10). Because a hot relay can no longer be flooded regardless of
// this number, we can fan out widely across the many well-behaved relays for
// throughput. RelayDiagnostics previously showed a blunt 24 caused
// rate-limited=1433 / "too many concurrent REQs"=1286; the adaptive cap
// targets exactly those relays instead of throttling everyone uniformly.
private const val DRAIN_CONCURRENCY = 48
// Sharded backbone sweep: instead of asking every popular relay for the
// same full author list (N× redundant), split the still-missing authors
@@ -337,7 +338,7 @@ object GrapeRankCommand {
val dead = hashSetOf<NormalizedRelayUrl>()
val filters =
mapOf(relay to shard.chunked(AUTHORS_PER_FILTER).map { Filter(kinds = graphKinds, authors = it) })
ctx.drain(filters, timeoutMs, diagnose, dead) to dead
ctx.drain(filters, timeoutMs, diagnose, dead, gatePerRelay = true) to dead
}
}
}.awaitAll()
@@ -363,7 +364,7 @@ object GrapeRankCommand {
val dead = hashSetOf<NormalizedRelayUrl>()
val filters =
live.associateWith { missing.chunked(AUTHORS_PER_FILTER).map { Filter(kinds = graphKinds, authors = it) } }
val events = ctx.drain(filters, timeoutMs, diagnose, dead)
val events = ctx.drain(filters, timeoutMs, diagnose, dead, gatePerRelay = true)
recordDead(dead)
relaysContacted += live
for ((relay, _) in events) liveRelays.add(relay)
@@ -417,7 +418,7 @@ object GrapeRankCommand {
.map { (batch, filters) ->
async {
val dead = hashSetOf<NormalizedRelayUrl>()
val events = ctx.drain(filters, timeoutMs, diagnose, dead)
val events = ctx.drain(filters, timeoutMs, diagnose, dead, gatePerRelay = true)
recordDead(dead)
Triple(batch, filters.keys, events)
}
@@ -473,6 +474,9 @@ object GrapeRankCommand {
if (ctx.relayDiagnostics.hadFeedback()) {
System.err.println("[graperank] relay feedback: ${ctx.relayDiagnostics.snapshot()}")
}
if (ctx.relayLimiter.hadThrottling()) {
System.err.println("[graperank] relay throttling: ${ctx.relayLimiter.snapshot()}")
}
} else {
// Offline: stream contact lists from the local store into the graph.
val loadStart = System.nanoTime()
@@ -530,6 +534,7 @@ object GrapeRankCommand {
"crawl_rounds" to rounds,
"relays_contacted" to relaysContactedCount,
"relay_feedback" to if (ctx.relayDiagnostics.hadFeedback()) ctx.relayDiagnostics.snapshot() else null,
"relay_throttling" to if (ctx.relayLimiter.hadThrottling()) ctx.relayLimiter.snapshot() else null,
"max_hop_reached" to (hopOf.values.maxOrNull() ?: 0),
"users_by_hop" to
hopOf.values
@@ -870,7 +875,7 @@ object GrapeRankCommand {
Filter(kinds = listOf(AdvertisedRelayListEvent.KIND), authors = chunk)
}
}
ctx.drain(filters, timeoutMs, diagnose)
ctx.drain(filters, timeoutMs, diagnose, gatePerRelay = true)
}
val discovery = relayListDiscoveryRelays(ctx)