From 73e68813f4604bf4e5a43d9d50b01004e162325b Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 3 Jul 2026 03:35:05 +0000 Subject: [PATCH] perf: shard PoolRequests lock per subscription MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The single global spin lock serialized every EVENT frame from every relay and measured negative scaling (4 concurrent relay consumers pushed 3.6M deliveries/s aggregate vs 11.1M for one thread alone). The lock now lives in RequestSubscriptionState — one per subscription — since all compound mutations are per-subId and different subs share no wire state. decideCommandLocked takes the state instance to avoid re-entering the non-reentrant lock; all-subs iterations lock one sub at a time; withLock is inline to keep the hot path allocation-free. DispatchStageBenchmark (PoolRequests-only, 1 -> 4 feeders): scaling flips from 0.33x to 3.4-7.3x across runs. PoolRequests concurrency, NostrClient and negentropy suites pass unchanged. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_018saXqYfAa3RvSJoDXK591R --- .../2026-07-02-nostrclient-receiver-perf.md | 17 ++ .../relay/client/pool/PoolRequests.kt | 153 ++++++++---------- .../client/reqs/RequestSubscriptionState.kt | 39 +++++ 3 files changed, 124 insertions(+), 85 deletions(-) diff --git a/quartz/plans/2026-07-02-nostrclient-receiver-perf.md b/quartz/plans/2026-07-02-nostrclient-receiver-perf.md index 452812148e..541663ca95 100644 --- a/quartz/plans/2026-07-02-nostrclient-receiver-perf.md +++ b/quartz/plans/2026-07-02-nostrclient-receiver-perf.md @@ -461,6 +461,23 @@ in ~45 min. To improve, in order: fine — that's TCP backpressure doing its job) and stream events to disk, never hold the set. +## Implemented optimizations (2026-07-03) + +### PoolRequests lock sharding (done) + +The global spin lock moved into `RequestSubscriptionState` — one lock per +subscription — since every compound mutation is scoped to one subId and +different subs share no wire state. `decideCommandLocked` now takes the state +instance so it never re-enters the (non-reentrant) lock, and the all-subs +iterations (connect/disconnect/cannot-connect) lock one sub at a time. +`withLock` is inline to keep the per-EVENT hot path allocation-free. + +Validation (DispatchStageBenchmark, PoolRequests-only, 4 feeder threads vs +1): scaling flipped from **negative** (11.1M → 3.6M deliveries/s, 0.33×) to +**positive** (3.4–7.3× across runs; absolute numbers on the shared CI box are +noisy, the scaling direction is consistent across every run). All +PoolRequests concurrency + NostrClient + negentropy suites pass. + ## Recommendations (in order of value/risk) 1. **Move Schnorr verification off the receiver coroutine** in the app's diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolRequests.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolRequests.kt index 8ed3822081..289d62300e 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolRequests.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolRequests.kt @@ -35,8 +35,6 @@ import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl import com.vitorpamplona.quartz.utils.cache.LargeCache import kotlinx.coroutines.flow.MutableStateFlow -import kotlin.concurrent.atomics.AtomicBoolean -import kotlin.concurrent.atomics.ExperimentalAtomicApi /** * Manages relay subscriptions for the entire pool in a way that only @@ -45,7 +43,6 @@ import kotlin.concurrent.atomics.ExperimentalAtomicApi * This code also awaits a subscription to come to EOSE since many relays * have through switching subs while they are processing the past. */ -@OptIn(ExperimentalAtomicApi::class) class PoolRequests { /** * Desired subs and listeners @@ -67,41 +64,29 @@ class PoolRequests { fun subState(subId: String): RequestSubscriptionState = relayState.getOrCreate(subId) { RequestSubscriptionState() } - /** - * Serializes every access to the subscription state machine - * ([RequestSubscriptionState]) and the "should I send a REQ?" decision. + /* + * Locking model: every compound access to a subscription's state machine + * ([RequestSubscriptionState]) — including the check-then-send decision in + * [decideCommandLocked] — runs inside THAT subscription's own lock + * ([RequestSubscriptionState.withLock]). * * A single subscription can span many relays, and each relay's - * socket-reader thread delivers messages into this class concurrently while - * the app thread adds/removes subscriptions — so the plain maps inside - * [RequestSubscriptionState] are written from several threads at once. That - * is both a memory hazard (concurrent map mutation) and, more importantly, - * a logic hazard: the check-then-send in [decideCommandLocked] must be - * atomic, otherwise two threads can both observe "no REQ in flight" and both - * send a REQ for the same sub id. + * socket-reader thread delivers messages into this class concurrently + * while the app thread adds/removes subscriptions — so the plain maps + * inside [RequestSubscriptionState] are written from several threads at + * once. That is both a memory hazard (concurrent map mutation) and a + * logic hazard: two threads must never both observe "no REQ in flight" + * and both send a REQ for the same sub id. * - * This is a tiny non-reentrant spin lock (the same [AtomicBoolean] primitive - * used by BasicRelayClient's connecting mutex): the critical sections are a - * handful of map operations, never any I/O. Listener callbacks and the - * actual socket sends are ALWAYS performed outside the lock — they re-enter - * this class through [onSent], so holding the lock across them would - * self-deadlock. + * The lock is per subscription, not global, because different subIds + * share no wire state: EVENT frames for different subs coming from + * different relay consumer threads must not serialize on each other (a + * global lock here measured negative scaling under 4 concurrent relay + * feeders). Listener callbacks and the actual socket sends are ALWAYS + * performed outside the lock — they re-enter this class through + * [onSent], so holding the lock across them would self-deadlock, and the + * lock is non-reentrant. Never hold two subscriptions' locks at once. */ - private val stateLock = AtomicBoolean(false) - - private inline fun withStateLock(block: () -> R): R { - while (stateLock.exchange(true)) { - // Another thread holds the lock. Spin-read until it looks free - // (test-and-test-and-set: cheaper on the cache line than hammering - // exchange) then retry the acquisition above. - while (stateLock.load()) { } - } - try { - return block() - } finally { - stateLock.store(false) - } - } /** * This is called when a sub is added or removed from this class and @@ -188,11 +173,9 @@ class PoolRequests { * When a connecting, updates the state of all subs */ fun onConnecting(url: NormalizedRelayUrl) { - // Change states to connecting. - withStateLock { - relayState.forEach { subId, state -> - state.connecting(url) - } + // Change states to connecting. One sub's lock at a time. + relayState.forEach { subId, state -> + state.withLock { state.connecting(url) } } } @@ -205,8 +188,8 @@ class PoolRequests { ) { when (cmd) { is ReqCmd -> { - withStateLock { - subState(cmd.subId).onOpenReq(relay, cmd.filters) + subState(cmd.subId).let { state -> + state.withLock { state.onOpenReq(relay, cmd.filters) } } desiredSubListeners.get(cmd.subId)?.onSubscriptionStarted( relay = relay.url, @@ -215,8 +198,8 @@ class PoolRequests { } is CloseCmd -> { - withStateLock { - subState(cmd.subId).onSubscriptionClosed(relay) + subState(cmd.subId).let { state -> + state.withLock { state.onSubscriptionClosed(relay) } } desiredSubListeners.get(cmd.subId)?.onSubscriptionClosed( relay = relay.url, @@ -236,11 +219,12 @@ class PoolRequests { is EventMessage -> { var isLive = false var forFilters: List? = null - withStateLock { - val state = relayState.get(msg.subId) - state?.onNewEvent(relay.url) - isLive = state?.currentState(relay.url) == ReqSubStatus.LIVE - forFilters = state?.lastKnownFilterStates(relay.url) + relayState.get(msg.subId)?.let { state -> + state.withLock { + state.onNewEvent(relay.url) + isLive = state.currentState(relay.url) == ReqSubStatus.LIVE + forFilters = state.lastKnownFilterStates(relay.url) + } } desiredSubListeners.get(msg.subId)?.onEvent( event = msg.event, @@ -253,14 +237,15 @@ class PoolRequests { is EoseMessage -> { var forFilters: List? = null val cmd = - withStateLock { - val state = relayState.get(msg.subId) - state?.onEose(relay.url) - forFilters = state?.lastKnownFilterStates(relay.url) - // Decide (and pre-mark) the resend while still holding the - // lock, so a concurrent subscribe/unsubscribe on the app - // thread can't also decide to send a REQ for this sub. - decideCommandLocked(msg.subId, relay.url) + relayState.get(msg.subId)?.let { state -> + state.withLock { + state.onEose(relay.url) + forFilters = state.lastKnownFilterStates(relay.url) + // Decide (and pre-mark) the resend while still holding the + // lock, so a concurrent subscribe/unsubscribe on the app + // thread can't also decide to send a REQ for this sub. + decideCommandLocked(state, msg.subId, relay.url) + } } desiredSubListeners.get(msg.subId)?.onEose( relay = relay.url, @@ -276,11 +261,12 @@ class PoolRequests { is ClosedMessage -> { var forFilters: List? = null val cmd = - withStateLock { - val state = relayState.get(msg.subId) - state?.onClosed(relay.url) - forFilters = state?.lastKnownFilterStates(relay.url) - decideCommandLocked(msg.subId, relay.url) + relayState.get(msg.subId)?.let { state -> + state.withLock { + state.onClosed(relay.url) + forFilters = state.lastKnownFilterStates(relay.url) + decideCommandLocked(state, msg.subId, relay.url) + } } desiredSubListeners.get(msg.subId)?.onClosed( message = msg.message, @@ -300,10 +286,8 @@ class PoolRequests { * When the relay disconnects */ fun onDisconnected(url: NormalizedRelayUrl) { - withStateLock { - relayState.forEach { subId, state -> - state.disconnected(url) - } + relayState.forEach { subId, state -> + state.withLock { state.disconnected(url) } } } @@ -326,20 +310,16 @@ class PoolRequests { url: NormalizedRelayUrl, errorMessage: String, ) { - // Snapshot the affected subs (and their last-known filters) under the - // lock, then notify listeners outside it. - val toNotify = - withStateLock { - val list = mutableListOf?>>() - relayState.forEach { subId, state -> - // These are all my subs.. need to figure out which relays have them - val subs = desiredSubs.get(subId) - if (subs != null && url in subs.keys) { - list.add(subId to state.lastKnownFilterStates(url)) - } - } - list + // Snapshot the affected subs (and their last-known filters) under each + // sub's own lock, then notify listeners outside it. + val toNotify = mutableListOf?>>() + relayState.forEach { subId, state -> + // These are all my subs.. need to figure out which relays have them + val subs = desiredSubs.get(subId) + if (subs != null && url in subs.keys) { + toNotify.add(subId to state.withLock { state.lastKnownFilterStates(url) }) } + } toNotify.forEach { (subId, forFilters) -> desiredSubListeners.get(subId)?.onCannotConnect( @@ -355,9 +335,10 @@ class PoolRequests { relaysToUpdate: Set, sync: (NormalizedRelayUrl, Command) -> Unit, ) { + val state = subState(subId) relaysToUpdate.forEach { relay -> - // Decide + pre-mark atomically under the lock, then send outside it. - val cmd = withStateLock { decideCommandLocked(subId, relay) } + // Decide + pre-mark atomically under the sub's lock, then send outside it. + val cmd = state.withLock { decideCommandLocked(state, subId, relay) } if (cmd != null) { sync(relay, cmd) } @@ -376,14 +357,16 @@ class PoolRequests { * CLOSE interleaves, an empty result that silently truncates a paged * download) — that is the bug this guards against. * - * MUST be called while holding [withStateLock]. + * MUST be called while holding [state]'s lock ([RequestSubscriptionState.withLock]); + * [state] must be [subId]'s state instance. Takes the instance instead of + * re-resolving it so it never re-enters the (non-reentrant) lock. */ private fun decideCommandLocked( + state: RequestSubscriptionState, subId: String, relay: NormalizedRelayUrl, ): Command? { - val state = relayState.get(subId) - val oldFilters = state?.currentFilters(relay) + val oldFilters = state.currentFilters(relay) val newFilters = desiredSubs.get(subId)?.get(relay) return when { @@ -401,12 +384,12 @@ class PoolRequests { // reply belongs to which REQ. The pending change is picked up // later by the EOSE handler, which runs this method again once the // sub reaches LIVE. - val current = state?.currentState(relay) + val current = state.currentState(relay) if (current == ReqSubStatus.SENT || current == ReqSubStatus.QUERYING_PAST) { null } else { // Pre-mark SENT + filters so a concurrent decider skips. - subState(subId).onOpenReq(relay, newFilters) + state.onOpenReq(relay, newFilters) ReqCmd(subId, newFilters) } } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/reqs/RequestSubscriptionState.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/reqs/RequestSubscriptionState.kt index f4c8bd8dd6..33b7238dec 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/reqs/RequestSubscriptionState.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/reqs/RequestSubscriptionState.kt @@ -21,12 +21,51 @@ package com.vitorpamplona.quartz.nip01Core.relay.client.reqs import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import kotlin.concurrent.atomics.AtomicBoolean +import kotlin.concurrent.atomics.ExperimentalAtomicApi /** * Manages the State of Subscriptions by logging states as the * subscription progresses. + * + * Thread-safety: the plain maps below are only touched inside [withLock]. + * The lock lives HERE — one per subscription — instead of a single global + * lock in PoolRequests, because every mutation is scoped to one subId: + * relays delivering EVENTs for *different* subscriptions have no shared + * state and must not serialize on each other. (A single global spin lock + * measured *negative* scaling: 4 relay consumer threads pushed less + * aggregate throughput through it than 1 — see + * quartz/plans/2026-07-02-nostrclient-receiver-perf.md.) */ +@OptIn(ExperimentalAtomicApi::class) class RequestSubscriptionState { + /** + * Tiny non-reentrant spin lock (same primitive as BasicRelayClient's + * connecting mutex). Critical sections are a handful of map operations, + * never I/O — callers MUST NOT re-enter and MUST NOT hold two + * subscriptions' locks at once (PoolRequests locks one sub at a time, + * including inside its all-subs iterations). + * + * `@PublishedApi internal` only because [withLock] is inline (this sits + * on the per-EVENT hot path; inlining avoids a closure allocation per + * message) — treat it as private. + */ + @PublishedApi + internal val lock = AtomicBoolean(false) + + inline fun withLock(block: () -> R): R { + while (lock.exchange(true)) { + // Test-and-test-and-set: spin-read until it looks free (cheaper + // on the cache line than hammering exchange), then retry above. + while (lock.load()) { } + } + try { + return block() + } finally { + lock.store(false) + } + } + // Logs the state of each channel to: // 1. inform when an event is received as live // 2. to block REQs being sent before finished (receiving an EOSE or Closed)