diff --git a/quartz/plans/2026-08-03-poolrequests-lock-contention.md b/quartz/plans/2026-08-03-poolrequests-lock-contention.md new file mode 100644 index 0000000000..5c14313e09 --- /dev/null +++ b/quartz/plans/2026-08-03-poolrequests-lock-contention.md @@ -0,0 +1,221 @@ +# PoolRequests subscription-state lock: what a suspending Mutex would cost + +**Date:** 2026-08-03 +**Status:** analysis + measurements; recommendation is NOT to use `Mutex` +**Harness:** `quartz/src/jvmTest/.../prodbench/LockDesignComparisonBenchmark.kt` +(`./gradlew :quartz:jvmTest --tests "*.LockDesignComparisonBenchmark"`) +**Context:** follows the Pixel 8 ANR (`anr_2026-08-03-12-55-26-256`) that replaced +`RequestSubscriptionState`'s busy-wait with a parking `PlatformLock`. + +## The question + +The parking lock removed the CPU burn, but it still **blocks** rather than suspends. +Relay consumers run as coroutines on `Dispatchers.IO` (limitedParallelism 64) and +production has ~191 live relays, so a blocked waiter occupies one of 64 dispatcher +threads. kotlinx's scheduler treats IO tasks as blocking and grows the pool when they +block — which is why the on-device fix moved CPU but left thread count flat (+7%). + +Would `kotlinx.coroutines.sync.Mutex` (a waiter *suspends*, freeing its thread) be +better? + +## What Mutex actually costs + +The lock sits at the bottom of a call chain that is **entirely non-suspend**: + +``` +OkHttpWebSocket: scope.launch { for (m in incomingMessages) out.onMessage(m) } <- coroutine + WebSocketListener.onMessage (non-suspend) + BasicRelayClient.MyWebsocketListener.onMessage + RelayPool.onIncomingMessage + NostrClient.onIncomingMessage (RelayConnectionListener) + PoolRequests.onIncomingMessage + RequestSubscriptionState.withLock <- the lock + SubscriptionListener.onEvent (non-suspend, fans out to the app) +``` + +`Mutex.lock()` is `suspend`, so every frame above it must become `suspend`. Measured +blast radius: + +| interface / API | count | +|---|---| +| `override fun onEvent(` | 59 | +| `override fun onEose(` | 52 | +| `override fun onIncomingMessage(` | 35 | +| `override fun onCannotConnect(` | 32 | +| `override fun onClosed(` | 25 | +| `override fun onDisconnected(` | 17 | +| `override fun onConnected(` | 16 | +| `override fun onSent(` | 14 | +| `override fun onConnecting(` | 11 | +| others (`onSubscriptionStarted`, …) | 1 | +| **total overrides to convert** | **262** | +| `.subscribe(` / `.unsubscribe(` call sites | **110** | + +Worse than the count: the *entry points* are not all coroutines. `INostrClient` is +non-suspend by design — `subscribe`, `unsubscribe`, `publish`, `syncFilters`, +`connect` — and is called from ViewModels, filter assemblers and Compose effects. +`PoolRequests.addOrUpdate` / `remove` / `sendToRelayIfChanged` reach the lock from +those paths. Making them suspend pushes coroutine scoping into every call site that +today just calls `subscribe(...)` synchronously. + +## The cheaper alternative: stripe the lock per relay + +**Enabling invariant (verified by reading all 11 `withLock` bodies):** *every* +critical section in `PoolRequests` is scoped to exactly one relay. Each one takes +`url` / `relay` / `relay.url` as its key: + +| line | body | key | +|---|---|---| +| 204 | `state.connecting(url)` | url | +| 218 | `state.onOpenReq(relay, cmd.filters)` | relay | +| 228 | `state.onSubscriptionClosed(relay)` | relay | +| 249 | `onNewEvent` / `currentState` / `lastKnownFilterStates` | relay.url | +| 267 | `onEose` + `decideCommandLocked` | relay.url | +| 296 | `onClosed` + `recordRefusalIfStructural` + `decideCommandLocked` | relay.url | +| 325 | `state.disconnected(url)` | url | +| 354 | `isStructurallyRefused` + `onOpenReq` | relay | +| 382 | `lastKnownFilterStates(url)` | url | +| 403 | `decideCommandLocked(state, subId, relay)` | relay | + +Every field in `RequestSubscriptionState` is a `Map` keyed by relay +(`subStates`, `filterStates`, `lastKnownFilterStates`, `refusedFilters`, +`refusalCounts`). The single cross-relay accessor, `currentFilters()` (no-arg, +returns the whole map), has **zero usages** — dead code, delete it. + +So the lock is only shared across relays as an artifact of `mutableMapOf` not being +thread-safe. One lock per `(subId, relay)` is semantically equivalent and drops +contention from ~191 threads to ~1–2 (that relay's consumer, plus the occasional app +thread in `sendToRelayIfChanged`). + +Blast radius: `RequestSubscriptionState` + `PoolRequests` only. Used by 4 production +files (`PoolRequests`, `RelayActiveRequestStates`, `RelayReqRefusals`) and 2 tests. +**No public API change.** + +## Measurements + +191 relay coroutines on a 64-thread limited-parallelism dispatcher, 3s windows. Each +iteration yields, modelling the real per-message suspension point of +`for (message in incomingMessages)` — without it the tight loops monopolise the +dispatcher and the harness measures itself, not the lock. + +`bystander` = a task that never touches the lock, on the same dispatcher. Its latency +answers "is the lock stealing dispatcher threads from unrelated work?" + +| subs | design | ops/s | bystander p50 | p99 | +|---|---|---|---|---| +| 1 | PER_SUB_BLOCKING (today) | 548,608 | 196.5µs | 770.9µs | +| 1 | PER_SUB_MUTEX | 190,672 | **1.9µs** | 30.4µs | +| 1 | **STRIPED** | **1,543,037** | 88.8µs | 237.6µs | +| 4 | PER_SUB_BLOCKING | 833,090 | 136.1µs | 507.6µs | +| 4 | PER_SUB_MUTEX | 523,440 | **3.1µs** | 26.0µs | +| 4 | **STRIPED** | **1,499,757** | 92.3µs | 337.5µs | +| 16 | PER_SUB_BLOCKING | 1,198,581 | 94.3µs | 444.3µs | +| 16 | PER_SUB_MUTEX | 749,487 | **5.0µs** | 16.5µs | +| 16 | **STRIPED** | **1,789,509** | 85.7µs | 219.8µs | + +Reading: + +- **STRIPED gives the most throughput** — 2.8× / 1.8× / 1.5× over today — and roughly + halves bystander p50. It removes the contention rather than tolerating it. +- **MUTEX is the slowest of the three** (0.35× / 0.63× / 0.63× of today): per-acquisition + suspend/resume is not free at these rates. +- **MUTEX does win bystander latency decisively** (1.9–5.0µs vs 85–196µs, and it collected + 626k vs 24k samples). But read that honestly — part of the win is that its coroutines + spend their time *suspended waiting for the mutex instead of doing work*. Better + thread-yielding is partly a symptom of lower throughput, not purely a win. + +## On-device result (SM-T220, playBenchmark, same account, n=3 per design) + +Cold-start burst, 75s window, exact per-thread CPU accounting from `/proc//stat`. +Warmups matched (~220-255 established sockets, ~200 relay reader threads each time). + +| design | DefaultDispatcher CPU (ms/75s) mean / median / range | GC ms | threads | +|---|---|---|---| +| spin lock (pre-fix) | 136,743 / 117,750 / 112,230–180,250 | 33,777 | 475 | +| parking, one lock per sub | 103,780 / 99,160 / 86,120–130,680 | 30,235 | 510 | +| **striped, per (sub, relay)** | **88,953 / 82,160 / 72,860–111,840** | **27,687** | 489 | + +- **Striped vs spin: −35% mean, −30% median CPU, −18% GC — and the ranges do not + overlap** (striped max 111,840 < spin min 112,230). That separation is what the + parking-only change could not show; its range overlapped the baseline heavily. +- **Striped vs parking: −14% mean / −17% median**, ranges still overlap — directional, + not conclusive at n=3. +- **Thread count is unchanged across all three** (475 / 510 / 489). Expected: the lock + is blocking, and kotlinx marks IO tasks blocking and grows the pool. Striping reduces + how *often* threads block, not the fact that they can. +- Not a false win from broken subscriptions: the feed rendered live notes 3-15 min old + with avatars and reaction counts, on 223 sockets / 200 relay reader threads. + +## Recommendation + +**Do the striping; do not convert to `Mutex`.** + +Striping attacks the cause — the lock is contended only because it is shared across +relays that touch disjoint keys. Once contention is ~0, the blocking-vs-suspending +question is moot: *a lock that is never contended never blocks a thread*. It also +happens to be the fastest option and is contained to two files, versus 262 overrides +and 110 call sites for a design that measures slower. + +`Mutex` only becomes the right answer if a future critical section must genuinely span +relays (or do I/O), which none does today. + +## The lock must not live inside a removable entry + +The obvious shape — `ConcurrentMap` where `PerRelayState` owns both +the fields *and* its lock — is **wrong as soon as entries can be removed**. +`connecting()` / `disconnected()` genuinely mean "forget this relay's wire state", so +they want removal, and then: + +``` +T1: getOrPut(R) -> stateA ; stateA.lock.lock() // in a critical section +T2: disconnected(R) -> remove(R) // stateA is now orphaned +T3: getOrPut(R) -> stateB (NEW lock) ; stateB.lock.lock() // acquires immediately + -> T1 and T3 are both "in" relay R's critical section, excluding nothing. +``` + +Two ways out: + +- **(A) never remove; null the fields.** Lock identity is stable, but entries + accumulate for every relay a sub has *ever* seen. Bounded but monotonic — and slow + monotonic growth in a long-lived process is precisely the class of bug this whole + investigation was about. +- **(B) stable stripe locks + removable state.** A fixed `Array(N) { PlatformLock() }` + indexed by `relay.hashCode()`, never mutated, so lock identity can't change; the + `ConcurrentMap` entries are then free to be added and removed with + exact `connecting`/`disconnected` semantics and no growth. + +**(B) is the recommendation.** With N = 32 and ~191 relays, ~6 relays share a stripe — +still a ~32x contention reduction versus today's single lock per sub, with none of the +lock-lifetime hazard. `ConcurrentMap.remove` (added 2026-08-03, with tests in +`ConcurrentCollectionsTest`) is what makes (B) possible. + +## Implementation sketch + +1. Delete the dead `currentFilters()` (no-arg). +2. In `RequestSubscriptionState`, replace the five relay-keyed maps with a single + `ConcurrentMap` holding the five fields — **no lock inside**. +3. Add a fixed `private val stripes = Array(32) { PlatformLock() }` and + `withLock(reference)` = `stripes[reference.hashCode().absoluteValue % 32].withLock { }`. + The array is never mutated, so lock identity is stable for the object's life. +4. `connecting()` / `disconnected()` become `map.remove(reference)` — exact current + semantics, no residue. +5. Update the 11 `withLock` call sites in `PoolRequests` to pass the relay. + `decideCommandLocked`'s "MUST hold the lock" contract becomes "MUST hold *that + relay's stripe*" — tighten the kdoc. +6. Extend `PoolRequestsRefusalTest` with concurrent multi-relay access on one subId, + and add a striped variant to `LockDesignComparisonBenchmark` to confirm the + measured win survives stripe collisions. + +## Risks + +- **Stripe collisions serialize unrelated relays.** With 32 stripes and 191 relays this + is ~6-way sharing; it is a contention *reduction*, not elimination. If a future + profile shows it mattering, raise N — it is a one-line change with no semantic effect. +- **Two relays on one stripe must never be locked simultaneously by one thread** — that + would self-deadlock (`PlatformLock` is reentrant on JVM/Apple, so same-thread + re-entry is survivable, but the invariant should be stated). `PoolRequests` already + locks one relay at a time. +- The state machine's atomicity comments are load-bearing; the per-relay scoping must + be re-verified against any new `withLock` body added between now and the change. +- Striping is NOT a fix for a critical section that does I/O or spans relays. If one is + ever added, this analysis must be redone. diff --git a/quartz/src/appleMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/PlatformLock.apple.kt b/quartz/src/appleMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/PlatformLock.apple.kt new file mode 100644 index 0000000000..c59d0ccecc --- /dev/null +++ b/quartz/src/appleMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/PlatformLock.apple.kt @@ -0,0 +1,36 @@ +/* + * 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.utils.concurrent + +import platform.Foundation.NSRecursiveLock + +// NSRecursiveLock parks contended waiters in the kernel, matching the +// ReentrantLock semantics of the jvmAndroid actual (and the same choice +// commons' KmpLock made for iOS). This must NOT be a spin lock: quartz's +// relay client runs here too, and spinning is what produced the Android +// ANR documented on the expect declaration. +actual class PlatformLock { + private val delegate = NSRecursiveLock() + + actual fun lock() = delegate.lock() + + actual fun unlock() = delegate.unlock() +} 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 eb8c220c7d..7d4b08315b 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 @@ -84,25 +84,26 @@ class PoolRequests( /* * 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]). + * [decideCommandLocked] — runs inside the stripe for THAT (subscription, relay) + * pair ([RequestSubscriptionState.withLock], which takes the relay). * - * 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 a - * logic hazard: two threads must never both observe "no REQ in flight" - * and both send a REQ for the same sub id. + * 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. Two threads must never both observe "no REQ in + * flight" and both send a REQ for the same (sub, relay). * - * 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. + * The lock is striped PER RELAY rather than per subscription: every critical + * section below touches only one relay's slice of the state, so EVENT frames + * arriving for the same subId from different relays no longer serialize on each + * other. With ~191 relays that previously made a single per-sub lock contended + * ~191 threads deep — the production ANR in + * `quartz/plans/2026-08-03-poolrequests-lock-contention.md`. + * + * Two rules this file must keep: + * - Never hold two relays' stripes (or two subscriptions') at once. Every loop + * below locks exactly one relay at a time. + * - Listener callbacks and socket sends are ALWAYS performed outside the lock — + * they re-enter this class through [onSent]. */ /** @@ -201,7 +202,7 @@ class PoolRequests( fun onConnecting(url: NormalizedRelayUrl) { // Change states to connecting. One sub's lock at a time. relayState.forEach { subId, state -> - state.withLock { state.connecting(url) } + state.withLock(url) { state.connecting(url) } } } @@ -215,7 +216,7 @@ class PoolRequests( when (cmd) { is ReqCmd -> { subState(cmd.subId).let { state -> - state.withLock { state.onOpenReq(relay, cmd.filters) } + state.withLock(relay) { state.onOpenReq(relay, cmd.filters) } } desiredSubListeners.get(cmd.subId)?.onSubscriptionStarted( relay = relay.url, @@ -225,7 +226,7 @@ class PoolRequests( is CloseCmd -> { subState(cmd.subId).let { state -> - state.withLock { state.onSubscriptionClosed(relay) } + state.withLock(relay) { state.onSubscriptionClosed(relay) } } desiredSubListeners.get(cmd.subId)?.onSubscriptionClosed( relay = relay.url, @@ -246,7 +247,7 @@ class PoolRequests( var isLive = false var forFilters: List? = null relayState.get(msg.subId)?.let { state -> - state.withLock { + state.withLock(relay.url) { state.onNewEvent(relay.url) isLive = state.currentState(relay.url) == ReqSubStatus.LIVE forFilters = state.lastKnownFilterStates(relay.url) @@ -264,7 +265,7 @@ class PoolRequests( var forFilters: List? = null val cmd = relayState.get(msg.subId)?.let { state -> - state.withLock { + state.withLock(relay.url) { state.onEose(relay.url) forFilters = state.lastKnownFilterStates(relay.url) // Decide (and pre-mark) the resend while still holding the @@ -293,7 +294,7 @@ class PoolRequests( var forFilters: List? = null val cmd = relayState.get(msg.subId)?.let { state -> - state.withLock { + state.withLock(relay.url) { state.onClosed(relay.url) forFilters = state.lastKnownFilterStates(relay.url) recordRefusalIfStructural(state, relay.url, msg.message, forFilters) @@ -322,7 +323,7 @@ class PoolRequests( */ fun onDisconnected(url: NormalizedRelayUrl) { relayState.forEach { subId, state -> - state.withLock { state.disconnected(url) } + state.withLock(url) { state.disconnected(url) } } } @@ -351,7 +352,7 @@ class PoolRequests( if (!filters.isNullOrEmpty()) { val send = subState(subId).let { state -> - state.withLock { + state.withLock(relay) { if (isStructurallyRefused(state, relay, filters)) { false } else { @@ -379,7 +380,7 @@ class PoolRequests( // 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.add(subId to state.withLock(url) { state.lastKnownFilterStates(url) }) } } @@ -400,7 +401,7 @@ class PoolRequests( val state = subState(subId) relaysToUpdate.forEach { relay -> // Decide + pre-mark atomically under the sub's lock, then send outside it. - val cmd = state.withLock { decideCommandLocked(state, subId, relay) } + val cmd = state.withLock(relay) { decideCommandLocked(state, subId, relay) } if (cmd != null) { sync(relay, cmd) } 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 4bfdd6aa2a..47b417241a 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,84 +21,116 @@ 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 +import com.vitorpamplona.quartz.utils.concurrent.ConcurrentMap +import com.vitorpamplona.quartz.utils.concurrent.PlatformLock /** * 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.) + * **Thread-safety: the lock is striped per reference (per relay), not per + * subscription.** Every operation below is scoped to a single [reference] — all state + * lives in [RelayState] objects held in a [ConcurrentMap] keyed by it — so two relays + * delivering EVENTs for the SAME subscription touch disjoint state and no longer + * serialize on each other. + * + * That mattered in production: one subId spans every relay it is subscribed on, so a + * single per-subscription lock was contended ~191 threads deep on the Pixel 8 that + * produced `anr_2026-08-03-12-55-26-256` (it was a *spin* lock then, which burned 6 of + * 9 cores — see [PlatformLock]). Striping removes the contention instead of making + * waiting cheaper: measured 1.5-2.8x the throughput of one lock per sub in + * `quartz/src/jvmTest/.../prodbench/LockDesignComparisonBenchmark.kt`. Full analysis, + * including why a suspending `Mutex` was rejected, is in + * `quartz/plans/2026-08-03-poolrequests-lock-contention.md`. + * + * The stripe array is allocated once and NEVER mutated, so a stripe's identity is + * stable for this object's whole life. That is load-bearing: if locks lived inside the + * per-relay values, a thread holding one while another thread dropped and re-created + * that entry would leave both "inside" the critical section excluding nothing. + * + * Callers MUST NOT hold two references' stripes at once ([PoolRequests] locks one relay + * at a time, including inside its all-subs iterations), and MUST keep socket sends and + * listener callbacks outside the critical section — they re-enter this class through + * `onSent`, and the point of a lock is to be held briefly. */ -@OptIn(ExperimentalAtomicApi::class) -class RequestSubscriptionState { +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). + * One reference's (relay's) slice of this subscription's state. Plain `var`s: every + * field is written and read under that reference's stripe by [PoolRequests]. * - * `@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. + * Note `RelayActiveRequestStates` uses this class WITHOUT locking. There the fields + * may be read stale — but the backing [ConcurrentMap] can no longer be structurally + * corrupted the way the plain `HashMap`s this replaced could. + */ + private class RelayState { + /** Null == no REQ state on this relay (fresh, or wiped by connecting/disconnected). */ + var status: ReqSubStatus? = null + + /** Filters of the REQ currently believed to be in flight. */ + var filters: List? = null + + /** + * Survives connect/disconnect so that if new events still arrive we can link + * them with the filters the relay was processing. + */ + var lastKnownFilters: List? = null + + /** + * Refused-filter memory. Unlike [status]/[filters] — per-connection wire state + * wiped by [connecting]/[disconnected] — this SURVIVES reconnects on purpose: a + * relay that structurally refuses a filter (a search-only relay CLOSING a plain + * kinds REQ, a relay that "does not accept REQs", "too many filters", …) refuses + * it again on every new socket, so [PoolRequests.syncState] replaying it each + * reconnect is pure waste. [refusalCount] accumulates repeated refusals of the + * same shape so a one-off (transient) close isn't mistaken for a structural one. + * Cleared on a successful REQ ([onEose]/[onNewEvent]) or when the caller observes + * the desired filter meaningfully changed. + */ + var refusedFilters: List? = null + + var refusalCount: Int = 0 + } + + private val states = ConcurrentMap() + + /** + * Fixed stripe array — allocated once, never mutated, so lock identity is stable. + * [STRIPE_COUNT] stripes over ~191 relays is roughly 6-way sharing: a ~32x + * contention reduction versus one lock per subscription. References colliding on a + * stripe merely serialize; correctness never depends on N. */ @PublishedApi - internal val lock = AtomicBoolean(false) + internal val stripes = Array(STRIPE_COUNT) { PlatformLock() } - 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()) { } - } + @PublishedApi + internal fun stripeFor(reference: T): PlatformLock = stripes[(reference.hashCode() and 0x7FFFFFFF) % STRIPE_COUNT] + + /** + * Runs [block] holding [reference]'s stripe. Inline so the per-EVENT hot path + * allocates no closure. + */ + inline fun withLock( + reference: T, + block: () -> R, + ): R { + val lock = stripeFor(reference) + lock.lock() try { return block() } finally { - lock.store(false) + lock.unlock() } } - // 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) - // - // If 2 happens, the relay might send multiple EOSEs in sequence - // for the same sub and we won't know which REQ was it for. - private val subStates = mutableMapOf() - private val filterStates = mutableMapOf>() + /** Read-only lookup — never creates an entry. */ + private fun peek(reference: T): RelayState? = states[reference] - /** - * This cache is used to make sure we know what the relay was processing - * before a close or disconnect so that if new events still arrive - * we can link them with the appropriate filters. - */ - private val lastKnownFilterStates = mutableMapOf>() + /** Write lookup — creates the entry on first use. */ + private fun mutable(reference: T): RelayState = states.getOrPut(reference) { RelayState() } - /** - * Refused-filter memory. Unlike [subStates]/[filterStates] above — per-connection - * wire state wiped by [connecting]/[disconnected] — this SURVIVES reconnects on - * purpose: a relay that structurally refuses a filter (a search-only relay CLOSING - * a plain kinds REQ, a relay that "does not accept REQs", "too many filters", …) - * refuses it again on every new socket, so [PoolRequests.syncState] replaying it - * each reconnect is pure waste. [refusalCounts] accumulates repeated refusals of - * the same shape so a one-off (transient) close isn't mistaken for a structural - * one. Cleared on a successful REQ ([onEose]/[onNewEvent]) or when the caller - * observes the desired filter meaningfully changed. - */ - private val refusedFilters = mutableMapOf>() - private val refusalCounts = mutableMapOf() + fun refusedFilters(reference: T) = peek(reference)?.refusedFilters - fun refusedFilters(reference: T) = refusedFilters[reference] - - fun refusalCount(reference: T) = refusalCounts[reference] ?: 0 + fun refusalCount(reference: T) = peek(reference)?.refusalCount ?: 0 /** * Records that [reference] refused [filters]. [sameAsLastRefusal] must be true when @@ -111,73 +143,91 @@ class RequestSubscriptionState { filters: List, sameAsLastRefusal: Boolean, ) { + val state = mutable(reference) if (sameAsLastRefusal) { - refusalCounts[reference] = refusalCount(reference) + 1 + state.refusalCount += 1 } else { - refusedFilters[reference] = filters - refusalCounts[reference] = 1 + state.refusedFilters = filters + state.refusalCount = 1 } } fun clearRefusal(reference: T) { - refusedFilters.remove(reference) - refusalCounts.remove(reference) + peek(reference)?.let { + it.refusedFilters = null + it.refusalCount = 0 + } } - fun currentFilters() = filterStates + fun currentFilters(reference: T) = peek(reference)?.filters - fun currentFilters(reference: T) = filterStates[reference] + fun lastKnownFilterStates(reference: T) = peek(reference)?.lastKnownFilters - fun lastKnownFilterStates(reference: T) = lastKnownFilterStates[reference] - - fun currentState(reference: T) = subStates[reference] + fun currentState(reference: T) = peek(reference)?.status fun onNewEvent(reference: T) { + val state = mutable(reference) // The relay is serving this REQ (it matched an event), so any past refusal // no longer applies — let it be tried freely again. - clearRefusal(reference) - if (subStates[reference] == ReqSubStatus.SENT) { - subStates[reference] = ReqSubStatus.QUERYING_PAST + state.refusedFilters = null + state.refusalCount = 0 + if (state.status == ReqSubStatus.SENT) { + state.status = ReqSubStatus.QUERYING_PAST } } fun onEose(reference: T) { + val state = mutable(reference) // Reaching EOSE means the relay accepted and finished the REQ; clear any refusal. - clearRefusal(reference) - subStates[reference] = ReqSubStatus.LIVE + state.refusedFilters = null + state.refusalCount = 0 + state.status = ReqSubStatus.LIVE } fun onClosed(reference: T) { - subStates[reference] = ReqSubStatus.CLOSED - // Closed messages are usually relays refusing to process a REQ - // This message keeps the state of filterStates intact to - // avoid sending the same filter, and getting immediately closed, - // over and over again. - - // filterStates.remove(reference) + // Closed messages are usually relays refusing to process a REQ. This keeps + // [RelayState.filters] intact to avoid sending the same filter, and getting + // immediately closed, over and over again. + mutable(reference).status = ReqSubStatus.CLOSED } fun onOpenReq( reference: T, filters: List, ) { - subStates[reference] = ReqSubStatus.SENT - filterStates[reference] = filters - lastKnownFilterStates[reference] = filters + val state = mutable(reference) + state.status = ReqSubStatus.SENT + state.filters = filters + state.lastKnownFilters = filters } fun onSubscriptionClosed(reference: T) { - subStates[reference] = ReqSubStatus.CLOSED - filterStates.remove(reference) + val state = mutable(reference) + state.status = ReqSubStatus.CLOSED + state.filters = null } fun connecting(reference: T) { - subStates.remove(reference) - filterStates.remove(reference) + // Wipes per-connection wire state only; lastKnownFilters and the refusal memory + // deliberately survive (see [RelayState]). + peek(reference)?.let { + it.status = null + it.filters = null + } } fun disconnected(reference: T) { - subStates.remove(reference) - filterStates.remove(reference) + peek(reference)?.let { + it.status = null + it.filters = null + } + } + + companion object { + /** + * Comfortably above the IO dispatcher's thread count (64 / 2) so collisions stay + * rare even when many workers are inside this class at once. + */ + const val STRIPE_COUNT = 32 } } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/LiveNegentropyIndex.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/LiveNegentropyIndex.kt index 0e0b134d63..9b39c70776 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/LiveNegentropyIndex.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/LiveNegentropyIndex.kt @@ -22,8 +22,7 @@ package com.vitorpamplona.quartz.nip77Negentropy import com.vitorpamplona.negentropy.storage.IStorage import com.vitorpamplona.quartz.nip01Core.store.IdAndTime -import kotlin.concurrent.atomics.AtomicBoolean -import kotlin.concurrent.atomics.ExperimentalAtomicApi +import com.vitorpamplona.quartz.utils.concurrent.PlatformLock /** * Always-current `(created_at, id)` index for NIP-77 negentropy — the @@ -54,13 +53,12 @@ import kotlin.concurrent.atomics.ExperimentalAtomicApi * sealed storage, so one snapshot backs any number of concurrent * sessions and stays valid even if the live index mutates after. * - * Thread-safety: all operations take a short spin lock (same pattern as - * `LiveEventStore`'s replay dedup). Mutations arrive from the store's - * single writer; snapshots from any REQ coroutine. + * Thread-safety: all operations take a short parking lock ([PlatformLock]). + * Mutations arrive from the store's single writer; snapshots from any REQ + * coroutine. */ -@OptIn(ExperimentalAtomicApi::class) class LiveNegentropyIndex { - private val lock = AtomicBoolean(false) + private val lock = PlatformLock() /** Sorted by (createdAt, id). Only touched under [locked]. */ private var entries = ArrayList() @@ -74,14 +72,17 @@ class LiveNegentropyIndex { private var cachedSnapshot: IStorage? = null private var cachedGeneration = -1L + // Parks rather than busy-waits. This lock's critical sections are FAR longer than a + // handful of map ops — [rebuild] sorts the whole entry list and snapshotting builds a + // storage — so a spinning waiter would burn a core for the duration of a sort while + // concurrent ingest threads pile up. Same defect that produced the client-side ANR + // documented on PlatformLock; see quartz/.../prodbench/SpinLockConvoyBenchmark.kt. private inline fun locked(block: () -> R): R { - while (lock.exchange(true)) { - while (lock.load()) { } - } + lock.lock() try { return block() } finally { - lock.store(false) + lock.unlock() } } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.kt index 47d3233c9b..f06d914583 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.kt @@ -61,6 +61,17 @@ expect class ConcurrentMap() { remap: (old: V, new: V) -> V, ): V + /** + * Removes [key] and returns the value it held, or null when absent. + * + * CAUTION: if the removed value owns a lock that callers acquire, removal + * breaks mutual exclusion — a thread holding the old value's lock and a + * thread that re-created the entry are no longer excluding each other. Keep + * such locks in a structure whose identity is stable (see + * `quartz/plans/2026-08-03-poolrequests-lock-contention.md`). + */ + fun remove(key: K): V? + fun size(): Int /** A point-in-time copy of the entries — safe to iterate without holding a lock. */ diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/PlatformLock.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/PlatformLock.kt new file mode 100644 index 0000000000..b43f69e5ec --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/PlatformLock.kt @@ -0,0 +1,75 @@ +/* + * 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.utils.concurrent + +/** + * A blocking mutual-exclusion lock whose waiters **park** (yield the core to the + * scheduler) instead of busy-waiting. + * + * Why this exists: `commonMain` has no `java.util.concurrent.locks.Lock`, so the + * relay client previously hand-rolled a spin lock over an `AtomicBoolean`. A + * busy-wait is only ever correct when the holder cannot be descheduled while + * holding the lock — and on Android that assumption is false. A production ANR on + * a Pixel 8 (`anr_2026-08-03-12-55-26-256`, Amethyst 1.13.1) caught the exact + * failure: the thread holding the lock was parked in `WaitingForGcToComplete` + * while **51 of 52** runnable relay-dispatch threads sat in the inlined spin loop, + * burning 596% CPU (6 of the phone's 9 cores) waiting for a holder that could not + * be scheduled to release it. The UI thread, needing to allocate, then waited on + * the same GC for over 5s and Android killed the frame with + * "Input dispatching timed out". + * + * Measured on `SpinLockConvoyBenchmark` (12-core dev machine — a phone is worse): + * the spin lock delivered 138M critical sections/s uncontended but only 1.2M/s + * with 52 contenders (0.86%, a 116x collapse), and an *unrelated* allocating + * thread's p90 latency went from 22µs to 10ms. Parking removes the CPU burn: a + * waiter costs one context switch instead of a whole core. + * + * Splits the same way [ConcurrentMap] does: + * - JVM / Android → `ReentrantLock` (real parking via `AbstractQueuedSynchronizer`). + * - Apple → `NSRecursiveLock`, which also parks. The relay client genuinely runs + * on iOS, so this must not spin — same choice commons' `KmpLock` already made. + * - Linux → a spin, since Kotlin/Native ships no parking lock and there is no + * Foundation; linuxX64 is a build/CI target, not a host for the many-relay + * workload. Swap in a pthread mutex if that changes. + * + * (`commons` has an equivalent `KmpLock`, but `commons` depends on `quartz` and not + * the reverse, so quartz cannot use it — keep the two in sync by hand.) + * + * Unlike the primitive it replaces, the JVM/Android actual is **reentrant**, so an + * accidental re-entry degrades into a no-op rather than a self-deadlock. Callers + * should still keep I/O and listener callbacks outside the critical section — that + * discipline is about holding the lock briefly, not about avoiding a hang. + */ +expect class PlatformLock() { + fun lock() + + fun unlock() +} + +/** Runs [block] holding [this]. Inline so hot paths allocate no closure. */ +inline fun PlatformLock.withLock(block: () -> R): R { + lock() + try { + return block() + } finally { + unlock() + } +} diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentCollectionsTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentCollectionsTest.kt index 208054f818..5f78d107f1 100644 --- a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentCollectionsTest.kt +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentCollectionsTest.kt @@ -71,6 +71,40 @@ class ConcurrentCollectionsTest { assertEquals(4, m["k"]) } + @Test + fun mapRemoveReturnsOldValueAndDeletes() { + val map = ConcurrentMap() + map["a"] = 1 + map["b"] = 2 + + assertEquals(1, map.remove("a")) + assertNull(map["a"]) + assertEquals(1, map.size()) + assertEquals(2, map["b"]) + } + + @Test + fun mapRemoveAbsentKeyIsNullAndNoOp() { + val map = ConcurrentMap() + map["b"] = 2 + + assertNull(map.remove("missing")) + assertEquals(1, map.size()) + assertEquals(2, map["b"]) + } + + @Test + fun mapRemoveThenGetOrPutRecreates() { + // The lifecycle PoolRequests needs: connecting()/disconnected() drop a relay's + // wire state, and the next REQ re-creates it. A stale value must never survive. + val map = ConcurrentMap() + map.getOrPut("relay") { 1 } + map.remove("relay") + + assertEquals(9, map.getOrPut("relay") { 9 }) + assertEquals(9, map["relay"]) + } + @Test fun mapSnapshotIsDetached() { val m = ConcurrentMap() diff --git a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.jvmAndroid.kt b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.jvmAndroid.kt index aaf40fdcb7..c68eb8da8a 100644 --- a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.jvmAndroid.kt +++ b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.jvmAndroid.kt @@ -49,6 +49,8 @@ actual class ConcurrentMap { remap: (old: V, new: V) -> V, ): V = map.merge(key, value) { old, new -> remap(old, new) }!! + actual fun remove(key: K): V? = map.remove(key) + actual fun size(): Int = map.size actual fun snapshot(): Map = HashMap(map) diff --git a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/utils/concurrent/PlatformLock.jvmAndroid.kt b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/utils/concurrent/PlatformLock.jvmAndroid.kt new file mode 100644 index 0000000000..58488489b6 --- /dev/null +++ b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/utils/concurrent/PlatformLock.jvmAndroid.kt @@ -0,0 +1,35 @@ +/* + * 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.utils.concurrent + +import java.util.concurrent.locks.ReentrantLock + +// ReentrantLock parks contended waiters through AbstractQueuedSynchronizer, so a +// thread waiting on a holder that lost its core (GC, preemption) costs one context +// switch rather than a spinning core. Non-fair on purpose: fairness would add a +// handoff per acquisition and the critical sections here are microseconds long. +actual class PlatformLock { + private val lock = ReentrantLock() + + actual fun lock() = lock.lock() + + actual fun unlock() = lock.unlock() +} diff --git a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/prodbench/LockDesignComparisonBenchmark.kt b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/prodbench/LockDesignComparisonBenchmark.kt new file mode 100644 index 0000000000..48e75ec97d --- /dev/null +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/prodbench/LockDesignComparisonBenchmark.kt @@ -0,0 +1,316 @@ +/* + * 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.prodbench + +import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.RequestSubscriptionState +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.launch +import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock +import org.junit.Test +import java.util.concurrent.ConcurrentHashMap +import java.util.concurrent.atomic.AtomicBoolean +import java.util.concurrent.atomic.AtomicLong +import java.util.concurrent.locks.ReentrantLock + +/** + * Compares the three candidate designs for `PoolRequests`' subscription-state lock, + * under the real production topology: R relay-consumer coroutines on a + * limited-parallelism dispatcher (`Dispatchers.IO` = 64) all delivering EVENTs for a + * small number of subscription ids. + * + * - **PER_SUB_BLOCKING** — today: one blocking lock per subId, shared by every relay + * that sub runs on. Waiters park, which fixed the CPU burn, but a parked waiter + * still OCCUPIES its dispatcher thread. + * - **PER_SUB_MUTEX** — kotlinx `Mutex`: a waiter *suspends* and releases its thread. + * Requires making the whole listener chain `suspend` (262 overrides, 110 call sites). + * - **STRIPED** — one lock per (subId, relay). Every critical section in PoolRequests + * is already scoped to a single relay, so this removes the contention outright and + * needs no API change. + * + * The metric that matters is NOT lock throughput — it is whether unrelated work can + * still get a thread while the lock is contended. `bystanderLatency` models that: a + * task that never touches the lock, submitted to the same dispatcher. + */ +class LockDesignComparisonBenchmark { + private val relays = 191 + private val dispatcherThreads = 64 + private val durationMs = 3000L + + /** One relay's slice of a subscription's state — all real fields are relay-keyed. */ + private class PerRelay { + val lock = ReentrantLock() + var status: Int = 0 + var filters: List? = null + var lastKnown: List? = null + } + + private class Striped { + val perRelay = ConcurrentHashMap() + + inline fun withLock( + relay: Int, + block: (PerRelay) -> R, + ): R { + val s = perRelay.computeIfAbsent(relay) { PerRelay() } + s.lock.lock() + try { + return block(s) + } finally { + s.lock.unlock() + } + } + } + + private class PerSubBlocking { + val lock = ReentrantLock() + val status = HashMap() + val filters = HashMap>() + } + + private class PerSubMutex { + val mutex = Mutex() + val status = HashMap() + val filters = HashMap>() + } + + private fun report( + name: String, + subs: Int, + ops: Long, + bystanderSamples: List, + ) { + val sorted = bystanderSamples.sorted() + val p50 = sorted[sorted.size / 2] / 1000.0 + val p99 = sorted[(sorted.size * 99) / 100] / 1000.0 + val max = sorted.last() / 1000.0 + println( + "%-18s subs=%-3d ops/s=%,10d bystander p50=%8.1fus p99=%9.1fus max=%9.1fus (n=%d)".format( + name, + subs, + ops * 1000 / durationMs, + p50, + p99, + max, + sorted.size, + ), + ) + } + + @Test + fun compareDesigns() { + if (System.getenv("PROD_RELAY_BENCH") == null && System.getProperty("prodRelayBench") == null) { + println("compareDesigns skipped. Run with -PprodRelayBench=1 to enable.") + return + } + println("relays=$relays dispatcherThreads=$dispatcherThreads window=${durationMs}ms") + println("bystander = a task that NEVER touches the lock, on the same dispatcher.") + println("Lower bystander latency = the lock is not stealing dispatcher threads.\n") + for (subs in listOf(1, 4, 16)) { + runPerSubBlocking(subs) + runPerSubMutex(subs) + runStriped(subs) + runRealStriped(subs) + println() + } + } + + private fun runPerSubBlocking(subs: Int) = + runBlocking { + @Suppress("DEPRECATION") + val dispatcher = Dispatchers.IO.limitedParallelism(dispatcherThreads) + val states = Array(subs) { PerSubBlocking() } + val stop = AtomicBoolean(false) + val ops = AtomicLong(0) + val bystander = ArrayList() + val jobs = + (0 until relays).map { relay -> + launch(dispatcher) { + var n = 0L + while (!stop.get()) { + val s = states[relay % subs] + s.lock.lock() + try { + s.status[relay] = 1 + s.filters[relay] = SAMPLE + s.status[relay] + } finally { + s.lock.unlock() + } + n++ + // Models `for (message in incomingMessages)`: every real + // iteration suspends, letting the dispatcher multiplex. + kotlinx.coroutines.yield() + } + ops.addAndGet(n) + } + } + val by = + launch(dispatcher) { + while (!stop.get()) { + val t = System.nanoTime() + kotlinx.coroutines.yield() + bystander.add(System.nanoTime() - t) + } + } + Thread.sleep(durationMs) + stop.set(true) + jobs.forEach { it.join() } + by.join() + report("PER_SUB_BLOCKING", subs, ops.get(), bystander.ifEmpty { listOf(0L) }) + } + + private fun runPerSubMutex(subs: Int) = + runBlocking { + @Suppress("DEPRECATION") + val dispatcher = Dispatchers.IO.limitedParallelism(dispatcherThreads) + val states = Array(subs) { PerSubMutex() } + val stop = AtomicBoolean(false) + val ops = AtomicLong(0) + val bystander = ArrayList() + val jobs = + (0 until relays).map { relay -> + launch(dispatcher) { + var n = 0L + while (!stop.get()) { + val s = states[relay % subs] + s.mutex.withLock { + s.status[relay] = 1 + s.filters[relay] = SAMPLE + s.status[relay] + } + n++ + // Models `for (message in incomingMessages)`: every real + // iteration suspends, letting the dispatcher multiplex. + kotlinx.coroutines.yield() + } + ops.addAndGet(n) + } + } + val by = + launch(dispatcher) { + while (!stop.get()) { + val t = System.nanoTime() + kotlinx.coroutines.yield() + bystander.add(System.nanoTime() - t) + } + } + Thread.sleep(durationMs) + stop.set(true) + jobs.forEach { it.join() } + by.join() + report("PER_SUB_MUTEX", subs, ops.get(), bystander.ifEmpty { listOf(0L) }) + } + + private fun runStriped(subs: Int) = + runBlocking { + @Suppress("DEPRECATION") + val dispatcher = Dispatchers.IO.limitedParallelism(dispatcherThreads) + val states = Array(subs) { Striped() } + val stop = AtomicBoolean(false) + val ops = AtomicLong(0) + val bystander = ArrayList() + val jobs = + (0 until relays).map { relay -> + launch(dispatcher) { + var n = 0L + while (!stop.get()) { + states[relay % subs].withLock(relay) { s -> + s.status = 1 + s.filters = SAMPLE + s.lastKnown = SAMPLE + } + n++ + // Models `for (message in incomingMessages)`: every real + // iteration suspends, letting the dispatcher multiplex. + kotlinx.coroutines.yield() + } + ops.addAndGet(n) + } + } + val by = + launch(dispatcher) { + while (!stop.get()) { + val t = System.nanoTime() + kotlinx.coroutines.yield() + bystander.add(System.nanoTime() - t) + } + } + Thread.sleep(durationMs) + stop.set(true) + jobs.forEach { it.join() } + by.join() + report("STRIPED", subs, ops.get(), bystander.ifEmpty { listOf(0L) }) + } + + /** + * Same topology, but driving the REAL shipped [RequestSubscriptionState] (striped per + * relay) instead of a prototype — so the measured win is a property of the code that + * ships, not of this file. + */ + private fun runRealStriped(subs: Int) = + runBlocking { + @Suppress("DEPRECATION") + val dispatcher = Dispatchers.IO.limitedParallelism(dispatcherThreads) + val states = Array(subs) { RequestSubscriptionState() } + val stop = AtomicBoolean(false) + val ops = AtomicLong(0) + val bystander = ArrayList() + val filters = listOf(Filter(kinds = listOf(1))) + val jobs = + (0 until relays).map { relay -> + launch(dispatcher) { + var n = 0L + val state = states[relay % subs] + while (!stop.get()) { + state.withLock(relay) { + state.onNewEvent(relay) + state.currentState(relay) + state.onOpenReq(relay, filters) + state.lastKnownFilterStates(relay) + } + n++ + kotlinx.coroutines.yield() + } + ops.addAndGet(n) + } + } + val by = + launch(dispatcher) { + while (!stop.get()) { + val t = System.nanoTime() + kotlinx.coroutines.yield() + bystander.add(System.nanoTime() - t) + } + } + Thread.sleep(durationMs) + stop.set(true) + jobs.forEach { it.join() } + by.join() + report("REAL_STRIPED", subs, ops.get(), bystander.ifEmpty { listOf(0L) }) + } + + companion object { + private val SAMPLE = listOf("kinds:1", "authors:abc") + } +} diff --git a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/prodbench/SpinLockConvoyBenchmark.kt b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/prodbench/SpinLockConvoyBenchmark.kt new file mode 100644 index 0000000000..c8b17e6d85 --- /dev/null +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/prodbench/SpinLockConvoyBenchmark.kt @@ -0,0 +1,242 @@ +/* + * 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.prodbench + +import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.RequestSubscriptionState +import org.junit.Assert.assertTrue +import org.junit.Test +import java.util.concurrent.CountDownLatch +import java.util.concurrent.atomic.AtomicBoolean +import java.util.concurrent.atomic.AtomicLong +import kotlin.concurrent.thread + +/** + * Reproduces the production ANR seen on a Pixel 8 (2026-08-03, anr_2026-08-03-12-55-26-256): + * 191 live relay sockets feed EVENT frames for the same handful of subscription ids, so + * dozens of DefaultDispatcher workers pile onto ONE [RequestSubscriptionState] busy-wait + * lock. In that trace 51 of 52 runnable workers sat in the inlined spin loop while the + * single lock holder was parked in `WaitingForGcToComplete`. + * + * Two things are measured: + * - [contendedThroughput]: aggregate critical sections/s as the contender count grows past + * the core count (the "negative scaling" the 2026-07-02 plan measured with 4 feeders, + * re-run at production fan-out). + * - [victimLatencyUnderSpin]: what an UNRELATED thread (stand-in for the UI thread) sees + * while the spinners run — this is the ANR mechanism, not the lock throughput. + */ +class SpinLockConvoyBenchmark { + private val cores = Runtime.getRuntime().availableProcessors() + + /** + * Every waiter targets ONE reference on purpose. The lock is striped per relay, so + * spreading threads over distinct relays would put them on different stripes and + * this benchmark would measure nothing — the point here is the primitive's behaviour + * when contention DOES land on a single stripe. + */ + private val hotRelay = 0 + + /** Mirrors the real critical section: a handful of map reads/writes, no I/O. */ + private fun criticalSection( + state: RequestSubscriptionState, + relay: Int, + ) { + state.onNewEvent(relay) + state.currentState(relay) + state.lastKnownFilterStates(relay) + } + + @Test + fun contendedThroughput() { + if (System.getenv("PROD_RELAY_BENCH") == null && System.getProperty("prodRelayBench") == null) { + println("contendedThroughput skipped. Run with -PprodRelayBench=1 to enable.") + return + } + println("cores = $cores") + println("threads | ops/s | vs 1 thread") + var baseline = 0.0 + for (threads in listOf(1, 2, 4, 8, 16, 32, 52)) { + val state = RequestSubscriptionState() + val stop = AtomicBoolean(false) + val ops = AtomicLong(0) + val start = CountDownLatch(1) + val workers = + (0 until threads).map { id -> + thread { + start.await() + var local = 0L + while (!stop.get()) { + state.withLock(hotRelay) { criticalSection(state, hotRelay) } + local++ + } + ops.addAndGet(local) + } + } + val t0 = System.nanoTime() + start.countDown() + Thread.sleep(2000) + stop.set(true) + workers.forEach { it.join() } + val secs = (System.nanoTime() - t0) / 1e9 + val rate = ops.get() / secs + if (threads == 1) baseline = rate + println("%7d | %9.0f | %.2fx".format(threads, rate, rate / baseline)) + } + } + + /** + * REGRESSION GUARD: contended waiters on [RequestSubscriptionState.withLock] must PARK, + * never busy-wait. Fails if the lock is ever turned back into a spin lock. + * + * This is the signature that identified the production ANR: in + * `anr_2026-08-03-12-55-26-256`, 37 of 52 runnable workers sat at one obfuscated line of + * `PoolRequests.onIncomingMessage` and 12 more at one line of `syncState$lambda$0` — all + * `state=R`, `sCount=0`, burning 596% CPU while the lock holder was stuck in + * `WaitingForGcToComplete`. With a spin lock this test observes ~40/40 waiters RUNNABLE; + * with a parking lock it observes 0. + */ + @Test + fun contendedWaitersParkInsteadOfSpinning() { + val spinners = 40 + val state = RequestSubscriptionState() + val stop = AtomicBoolean(false) + val start = CountDownLatch(1) + + val holder = + thread(name = "holder") { + start.await() + while (!stop.get()) { + state.withLock(hotRelay) { + val t = System.nanoTime() + while (System.nanoTime() - t < 5_000_000) { /* hold 5ms */ } + } + } + } + val threads = + (0 until spinners).map { id -> + thread(name = "spinner-$id") { + start.await() + while (!stop.get()) { + state.withLock(hotRelay) { criticalSection(state, hotRelay) } + } + } + } + start.countDown() + Thread.sleep(300) + + // Sample every spinner's stack, exactly like an ANR dump would. + val points = HashMap() + var runnable = 0 + threads.forEach { t -> + if (t.state == Thread.State.RUNNABLE) runnable++ + val top = t.stackTrace.firstOrNull { it.className.contains("SpinLockConvoyBenchmark") || it.className.contains("RequestSubscriptionState") } + if (top != null) { + val key = "${top.className.substringAfterLast('.')}.${top.methodName}:${top.lineNumber}" + points[key] = (points[key] ?: 0) + 1 + } + } + stop.set(true) + threads.forEach { it.join() } + holder.join() + + println("sampled $spinners waiters: RUNNABLE=$runnable, distinct program points=${points.size}") + points.entries.sortedByDescending { it.value }.forEach { println(" ${it.value}x ${it.key}") } + + // A parked waiter reports WAITING, so the parking lock measures ~0 here while a + // spin lock measures ~100%. The line sits at 50% deliberately: it separates the + // two cases by a mile and leaves headroom on a loaded CI box, where a few waiters + // can legitimately be mid-acquire when we sample. + assertTrue( + "Expected contended waiters to park, but $runnable/$spinners were RUNNABLE — " + + "RequestSubscriptionState.withLock looks like it is busy-waiting again. " + + "See PlatformLock's kdoc: this is what caused anr_2026-08-03-12-55-26-256.", + runnable <= spinners / 2, + ) + } + + @Test + fun victimLatencyUnderSpin() { + if (System.getenv("PROD_RELAY_BENCH") == null && System.getProperty("prodRelayBench") == null) { + println("victimLatencyUnderSpin skipped. Run with -PprodRelayBench=1 to enable.") + return + } + // One holder that briefly stalls inside the critical section (in production: a GC + // pause, or simply being descheduled). Everyone else spins. + val spinners = 52 + val state = RequestSubscriptionState() + val stop = AtomicBoolean(false) + + fun measureVictim(label: String) { + // The "UI thread": allocates and measures its own scheduling latency. + val samples = ArrayList() + repeat(200) { + val t = System.nanoTime() + // trivial allocation work, like a Compose semantics traversal step + val junk = ArrayList(64) + repeat(64) { i -> junk.add("node$i") } + Thread.yield() + samples.add(System.nanoTime() - t) + } + samples.sort() + println( + "%-22s p50=%6.0fus p90=%7.0fus p99=%8.0fus max=%8.0fus".format( + label, + samples[samples.size / 2] / 1000.0, + samples[(samples.size * 90) / 100] / 1000.0, + samples[(samples.size * 99) / 100] / 1000.0, + samples.last() / 1000.0, + ), + ) + } + + measureVictim("idle (no spinners)") + + val start = CountDownLatch(1) + val holder = + thread { + start.await() + while (!stop.get()) { + state.withLock(hotRelay) { + // Simulate the holder losing its core / waiting on GC mid-section. + val t = System.nanoTime() + while (System.nanoTime() - t < 2_000_000) { /* 2ms stall */ } + } + Thread.sleep(1) + } + } + val threads = + (0 until spinners).map { id -> + thread { + start.await() + while (!stop.get()) { + state.withLock(hotRelay) { criticalSection(state, hotRelay) } + } + } + } + start.countDown() + Thread.sleep(500) + measureVictim("$spinners spinners") + stop.set(true) + threads.forEach { it.join() } + holder.join() + + measureVictim("after (no spinners)") + } +} diff --git a/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/PlatformLock.linux.kt b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/PlatformLock.linux.kt new file mode 100644 index 0000000000..5f3c4fdf46 --- /dev/null +++ b/quartz/src/linuxMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/PlatformLock.linux.kt @@ -0,0 +1,44 @@ +/* + * 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.utils.concurrent + +import kotlin.concurrent.atomics.AtomicBoolean +import kotlin.concurrent.atomics.ExperimentalAtomicApi + +// Linux/JVM-less target: Kotlin/Native's stdlib ships no parking lock and there is +// no Foundation here, so this keeps a test-and-test-and-set spin. Correct but not +// scalable. Acceptable ONLY because linuxX64 is a build/CI target for quartz, not a +// host for the many-relay client workload whose contention motivated the parking +// actuals on jvmAndroid and Apple. If that ever changes, swap in a pthread mutex. +@OptIn(ExperimentalAtomicApi::class) +actual class PlatformLock { + private val held = AtomicBoolean(false) + + actual fun lock() { + while (held.exchange(true)) { + while (held.load()) { } + } + } + + actual fun unlock() { + held.store(false) + } +} diff --git a/quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.native.kt b/quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.native.kt index de4ee540d8..8805e8d5a7 100644 --- a/quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.native.kt +++ b/quartz/src/nativeMain/kotlin/com/vitorpamplona/quartz/utils/concurrent/ConcurrentMap.native.kt @@ -74,6 +74,16 @@ actual class ConcurrentMap { } } + actual fun remove(key: K): V? { + while (true) { + val cur = ref.load() + val old = cur[key] ?: return null + val copy = HashMap(cur) + copy.remove(key) + if (ref.compareAndSet(cur, copy)) return old + } + } + actual fun size(): Int = ref.load().size actual fun snapshot(): Map = HashMap(ref.load())