From ca89b30fbe19cbfeee7886cee15dbeb753adf35f Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 19 Aug 2026 21:27:15 +0000 Subject: [PATCH 1/2] fix(quartz): release pool request/count state rows when a sub is removed MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit PoolRequests.remove() dropped the sub from `desiredSubs` and its listener, but left its RequestSubscriptionState in `relayState` forever. Since onConnecting/onDisconnected/onCannotConnect scan that whole map on every relay lifecycle event, each connect was O(all subs ever created) rather than O(live subs). A long-running pool driving many relays therefore spends more and more dispatcher time in the connect path, starving everything else sharing the client (reported downstream in NosFabrica/vespa-relay#154, where the monitor plane completed zero passes in 3h while ~95% of dispatcher CPU sat in the registry scan). PoolCounts had the same leak, plus a cross-registry one: NostrClient fans every frame out to both registries, so each REQ CLOSE materialized a permanent COUNT row (and each COUNT id a permanent REQ row via sendToRelayIfChanged). The registry deliberately models the *believed relay state*, separate from the desired state, so a row must outlive `remove()` — the CLOSE frame is decided from the filters it still holds. Dropping the row inside `remove()` looks like a one-line fix but silently stops CLOSE being sent, leaking the subscription on the relay instead. So the row is released after the CLOSE has been handed to the socket, and the write paths no longer resurrect it: - remove(): keeps the row (comment says why); drops it when the id isn't ours. - sendToRelayIfChanged(): never creates a row; releases it once the sub is no longer desired and every affected relay has been told. - onSent(): non-creating lookups. Both send paths pre-mark the row under the lock before the frame leaves, so a null means the sub is already gone. - PoolCounts: same, via a liveState() helper gated on the query still existing. While a sub is still desired nothing changes — its row survives connect, disconnect and filter changes, so late events still link to the filters the relay was running. Once removed, remove() has already dropped the listener in the same call, so the surviving row had no consumer left anyway. Adds INostrClient.registrySizes() so a host can alarm on registry growth; registry size multiplies the cost of every relay lifecycle event, so a leak shows up as connect-path CPU long before it shows up as memory. Tests cover both directions: the registry must stay bounded across churn, and CLOSE must still be sent for every removal (the guard that catches the naive one-line fix). Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01WHsz9fXDYysZcKZXhff4XY --- .../nip01Core/relay/client/INostrClient.kt | 18 + .../nip01Core/relay/client/NostrClient.kt | 6 + .../nip01Core/relay/client/pool/PoolCounts.kt | 38 +- .../relay/client/pool/PoolRequests.kt | 45 ++- .../pool/PoolRequestsRegistryLeakTest.kt | 345 ++++++++++++++++++ 5 files changed, 443 insertions(+), 9 deletions(-) create mode 100644 quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolRequestsRegistryLeakTest.kt diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt index bbed0d924e..0596190836 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/INostrClient.kt @@ -143,6 +143,18 @@ interface INostrClient : AutoCloseable { fun activeCounts(url: NormalizedRelayUrl): Map> + /** + * Size of the pool's REQ / COUNT state registries — the maps every relay + * connect, disconnect and connection failure scans in full. + * + * These track LIVE subscriptions, so a healthy client keeps them near the number of + * subscriptions it believes it holds. A host that drives many relays from one client + * should alarm on unbounded growth here: registry size multiplies the cost of every + * relay lifecycle event, so a leak shows up as connect-path CPU long before it shows + * up as memory. + */ + fun registrySizes(): RegistrySizes = RegistrySizes(0, 0) + fun activeOutboxCache(url: NormalizedRelayUrl): Set /** The events still pending delivery to [url] (full events, not just ids). */ @@ -205,3 +217,9 @@ class EmptyNostrClient : INostrClient { override fun close() {} } + +/** @see INostrClient.registrySizes */ +data class RegistrySizes( + val liveRequests: Int, + val liveCounts: Int, +) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt index 30da4b8992..95c5cda287 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt @@ -458,6 +458,12 @@ class NostrClient( override fun activeCounts(url: NormalizedRelayUrl): Map> = activeCounts.activeFiltersFor(url) + override fun registrySizes() = + RegistrySizes( + liveRequests = activeRequests.activeSubscriptionCount(), + liveCounts = activeCounts.activeQueryCount(), + ) + override fun activeOutboxCache(url: NormalizedRelayUrl): Set = eventOutbox.activeOutboxCacheFor(url) override fun activeOutboxEvents(url: NormalizedRelayUrl): List = eventOutbox.activeOutboxEventsFor(url) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolCounts.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolCounts.kt index 97c48dc5fa..34aefc19cd 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolCounts.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolCounts.kt @@ -38,9 +38,26 @@ class PoolCounts { private var queries = LargeCache>>() val relays = MutableStateFlow(setOf()) + /** + * **Invariant: LIVE queries only** -- a row exists only while [queries] still wants + * the id. [onConnecting] and [onDisconnected] scan this whole map on every relay + * lifecycle event, so a registry that accumulated dead ids would make each connect + * O(all queries ever made). See the matching note in [PoolRequests]. + */ private val relayState = LargeCache>() - fun subState(subId: String): CountQueryState = relayState.getOrCreate(subId) { CountQueryState() } + private fun subState(subId: String): CountQueryState = relayState.getOrCreate(subId) { CountQueryState() } + + /** + * State row for a query that is still desired, creating it on first use. Returns + * null once the query is removed, so a frame still on the wire (or a CLOSE ack for + * a REQ id, which NostrClient also routes through this registry) can never + * resurrect a dead row. + */ + private fun liveState(queryId: String): CountQueryState? = if (queries.containsKey(queryId)) subState(queryId) else null + + /** Live query count. Exposed so a host can alarm on registry growth. */ + fun activeQueryCount(): Int = relayState.size() private fun updateRelays() { val myRelays = mutableSetOf() @@ -85,8 +102,12 @@ class PoolCounts { val oldRelays = queries.get(queryId)?.keys ?: emptySet() queries.remove(queryId) updateRelays() + // The state row survives until [sendToRelayIfChanged] flushes the CLOSE -- + // that decision compares against the filters this relay still has in flight. oldRelays } else { + // Not ours (a REQ id, or a double remove): drop any dead row. + relayState.remove(queryId) emptySet() } @@ -118,8 +139,10 @@ class PoolCounts { cmd: Command, ) { when (cmd) { - is CountCmd -> subState(cmd.queryId).onQuery(relay, cmd.filters) - is CloseCmd -> subState(cmd.subId).onCloseQuery(relay) + is CountCmd -> liveState(cmd.queryId)?.onQuery(relay, cmd.filters) + // NostrClient routes every CLOSE through both registries, so the vast + // majority of these ids are REQ subs this class has never heard of. + is CloseCmd -> liveState(cmd.subId)?.onCloseQuery(relay) } } @@ -129,14 +152,14 @@ class PoolCounts { ) { when (msg) { is CountMessage -> { - subState(msg.queryId).onCountReply(relay.url) + liveState(msg.queryId)?.onCountReply(relay.url) sendToRelayIfChanged(msg.queryId, relay.url) { cmd -> relay.sendOrConnectAndSync(cmd) } } is ClosedMessage -> { - subState(msg.subId).onClosed(relay.url) + liveState(msg.subId)?.onClosed(relay.url) sendToRelayIfChanged(msg.subId, relay.url) { cmd -> relay.sendOrConnectAndSync(cmd) } @@ -167,6 +190,11 @@ class PoolCounts { } } } + + // Query is gone and every relay has been told: drop the row (see [relayState]). + if (!queries.containsKey(queryId)) { + relayState.remove(queryId) + } } fun sendToRelayIfChanged( 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 7c87d14712..11d8c38e77 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 @@ -76,10 +76,22 @@ class PoolRequests( * * This is important to link responses with which filter was a sub replying to and * to figure out if we need to update the relay if a REQ has changed. + * + * **Invariant: this holds LIVE subscriptions only.** A row exists between + * [addOrUpdate] and the [sendToRelayIfChanged] that flushes the CLOSE for a + * [remove]d sub, and nothing else may create one (see [onSent]). That invariant is + * what bounds [onConnecting] / [onDisconnected] / [onCannotConnect], which scan the + * whole map on every relay lifecycle event: a registry that instead accumulated + * every sub id ever used makes each connect O(all subs ever created), and on a + * server-side pool driving thousands of relays those scans saturate the dispatcher + * and starve every other plane sharing the client. */ private val relayState = LargeCache>() - fun subState(subId: String): RequestSubscriptionState = relayState.getOrCreate(subId) { RequestSubscriptionState() } + private fun subState(subId: String): RequestSubscriptionState = relayState.getOrCreate(subId) { RequestSubscriptionState() } + + /** Live subscription count. Exposed so a host can alarm on registry growth. */ + fun activeSubscriptionCount(): Int = relayState.size() /* * Locking model: every compound access to a subscription's state machine @@ -184,9 +196,17 @@ class PoolRequests( desiredSubListeners.remove(subId) // update relays for pool updateRelays() + // NOTE: the state row deliberately survives until [sendToRelayIfChanged] + // flushes the CLOSE -- that decision reads the filters this relay still + // believes are in flight. Dropping the row here makes the CLOSE frame + // disappear (the sub then leaks on the RELAY instead). // return all affected relays oldRelays } else { + // Nothing desired under this id, so any row left behind is dead: a double + // unsubscribe, or the same id being removed from the sibling registry + // (NostrClient.unsubscribe hits both this and PoolCounts). + relayState.remove(subId) emptySet() } @@ -215,7 +235,12 @@ class PoolRequests( ) { when (cmd) { is ReqCmd -> { - subState(cmd.subId).let { state -> + // Non-creating on purpose: both send paths ([decideCommandLocked] and + // [syncState]) pre-mark the row under the lock before the frame leaves, + // so a row always exists for a REQ we actually sent. A null here means + // the sub was unsubscribed while the frame was in flight -- recreating + // the row would resurrect a dead subscription in the registry. + relayState.get(cmd.subId)?.let { state -> state.withLock(relay) { state.onOpenReq(relay, cmd.filters) } } desiredSubListeners.get(cmd.subId)?.onSubscriptionStarted( @@ -225,7 +250,7 @@ class PoolRequests( } is CloseCmd -> { - subState(cmd.subId).let { state -> + relayState.get(cmd.subId)?.let { state -> state.withLock(relay) { state.onSubscriptionClosed(relay) } } desiredSubListeners.get(cmd.subId)?.onSubscriptionClosed( @@ -398,7 +423,11 @@ class PoolRequests( relaysToUpdate: Set, sync: (NormalizedRelayUrl, Command) -> Unit, ) { - val state = subState(subId) + // Never create a row here. NostrClient.unsubscribe calls this on BOTH registries + // with the same id, so a COUNT id would otherwise materialize a permanent + // (and never-scanned-past) row in the REQ registry, and vice-versa. + val state = relayState.get(subId) ?: return + relaysToUpdate.forEach { relay -> // Decide + pre-mark atomically under the sub's lock, then send outside it. val cmd = state.withLock(relay) { decideCommandLocked(state, subId, relay) } @@ -406,6 +435,14 @@ class PoolRequests( sync(relay, cmd) } } + + // The sub is gone and every affected relay has now been told. Drop the row so + // the registry keeps holding live subscriptions only. Frames that were already + // on the wire when the CLOSE went out land in [onIncomingMessage], find no row + // and no listener, and correctly no-op. + if (!desiredSubs.containsKey(subId)) { + relayState.remove(subId) + } } /** diff --git a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolRequestsRegistryLeakTest.kt b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolRequestsRegistryLeakTest.kt new file mode 100644 index 0000000000..a3c7ed9c12 --- /dev/null +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolRequestsRegistryLeakTest.kt @@ -0,0 +1,345 @@ +/* + * 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.client.pool + +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener +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.EoseMessage +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.EventMessage +import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.CloseCmd +import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.Command +import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.CountCmd +import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.ReqCmd +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import kotlinx.coroutines.test.runTest +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertTrue + +/** + * The request registry must hold LIVE subscriptions only. + * + * [PoolRequests.onConnecting], [PoolRequests.onDisconnected] and + * [PoolRequests.onCannotConnect] scan the whole registry on every relay lifecycle + * event. `remove()` used to drop the sub from `desiredSubs` but leave its + * [com.vitorpamplona.quartz.nip01Core.relay.client.reqs.RequestSubscriptionState] + * behind forever, so the registry grew with every one-shot fetch and each connect + * became O(all subs ever created) — reported downstream as a pool whose monitor plane + * never completed a pass because the dispatcher was pinned scanning dead rows + * (NosFabrica/vespa-relay#154). + * + * The naive fix — dropping the row inside `remove()` — is worse than the leak: the + * CLOSE frame is decided from that row, so it silently stops being sent and the + * subscription leaks on the *relay* instead. Hence [closesAreStillSentForEveryRemoval]. + */ +class PoolRequestsRegistryLeakTest { + private val relayA = NormalizedRelayUrl("wss://a.example/") + private val relayB = NormalizedRelayUrl("wss://b.example/") + + private fun filters() = listOf(Filter(kinds = listOf(1))) + + private fun event() = + Event( + id = "00".repeat(32), + pubKey = "11".repeat(32), + createdAt = 1, + kind = 1, + tags = emptyArray(), + content = "hi", + sig = "22".repeat(64), + ) + + private class FakeRelayClient( + override val url: NormalizedRelayUrl, + ) : IRelayClient { + override fun connect() = Unit + + override fun needsToReconnect() = false + + override fun connectAndSyncFiltersIfDisconnected(ignoreRetryDelays: Boolean) = Unit + + override fun isConnected() = true + + override fun disconnect() = Unit + + override fun sendIfConnected(cmd: Command) = Unit + + override fun sendOrConnectAndSync(cmd: Command) = Unit + } + + /** Mimics NostrClient.unsubscribe: remove from BOTH registries, then flush. */ + private fun unsubscribe( + reqs: PoolRequests, + counts: PoolCounts, + subId: String, + sent: MutableList, + ) { + val reqRelays = reqs.remove(subId) + val countRelays = counts.remove(subId) + reqs.sendToRelayIfChanged(subId, reqRelays) { r, cmd -> + sent.add(cmd) + // NostrClient.onSent fans every frame out to both registries. + reqs.onSent(r, cmd) + counts.onSent(r, cmd) + } + counts.sendToRelayIfChanged(subId, countRelays) { r, cmd -> + sent.add(cmd) + reqs.onSent(r, cmd) + counts.onSent(r, cmd) + } + } + + private fun subscribe( + reqs: PoolRequests, + counts: PoolCounts, + subId: String, + relays: Set, + sent: MutableList = mutableListOf(), + listener: SubscriptionListener? = null, + ) { + val toUpdate = reqs.addOrUpdate(subId, relays.associateWith { filters() }, listener) + reqs.sendToRelayIfChanged(subId, toUpdate) { r, cmd -> + sent.add(cmd) + reqs.onSent(r, cmd) + counts.onSent(r, cmd) + } + } + + @Test + fun registryHoldsLiveSubscriptionsOnly() { + val reqs = PoolRequests() + val counts = PoolCounts() + val sent = mutableListOf() + + repeat(1000) { i -> + val subId = "one-shot-$i" + subscribe(reqs, counts, subId, setOf(relayA, relayB), sent) + assertEquals(1, reqs.activeSubscriptionCount(), "sub $i must be live while subscribed") + unsubscribe(reqs, counts, subId, sent) + } + + assertEquals(0, reqs.activeSubscriptionCount(), "REQ registry must not retain removed subscriptions") + assertEquals(0, counts.activeQueryCount(), "CLOSEs for REQ subs must not materialize COUNT rows") + } + + @Test + fun closesAreStillSentForEveryRemoval() { + val reqs = PoolRequests() + val counts = PoolCounts() + val sent = mutableListOf() + + repeat(100) { i -> + val subId = "one-shot-$i" + subscribe(reqs, counts, subId, setOf(relayA, relayB), sent) + unsubscribe(reqs, counts, subId, sent) + } + + assertEquals(200, sent.count { it is ReqCmd }, "one REQ per (sub, relay)") + assertEquals(200, sent.count { it is CloseCmd }, "one CLOSE per (sub, relay) — the relay must be told") + assertEquals(0, sent.count { it is CountCmd }) + } + + @Test + fun countQueriesAreAlsoReleased() { + val reqs = PoolRequests() + val counts = PoolCounts() + val sent = mutableListOf() + + repeat(500) { i -> + val queryId = "count-$i" + val toUpdate = counts.addOrUpdate(queryId, mapOf(relayA to filters())) + counts.sendToRelayIfChanged(queryId, toUpdate) { r, cmd -> + sent.add(cmd) + reqs.onSent(r, cmd) + counts.onSent(r, cmd) + } + assertEquals(1, counts.activeQueryCount()) + unsubscribe(reqs, counts, queryId, sent) + } + + // NB: a COUNT still awaiting its reply is deliberately NOT closed (the reply + // ends the query on its own) — that guard is untouched here; what must change + // is that its state row is released. + assertEquals(500, sent.count { it is CountCmd }, "one COUNT per query") + assertEquals(0, counts.activeQueryCount(), "COUNT registry must not retain removed queries") + assertEquals(0, reqs.activeSubscriptionCount(), "CLOSEs for COUNT ids must not materialize REQ rows") + } + + /** Frames already on the wire when the CLOSE went out must no-op, not resurrect. */ + @Test + fun lateFramesAfterCloseDoNotResurrectTheRegistry() = + runTest { + val reqs = PoolRequests() + val counts = PoolCounts() + val sent = mutableListOf() + val relay = FakeRelayClient(relayA) + + var callbacks = 0 + val listener = + object : SubscriptionListener { + override fun onEose( + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + callbacks++ + } + + override fun onClosed( + message: String, + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + callbacks++ + } + } + + subscribe(reqs, counts, "sub", setOf(relayA), sent, listener) + unsubscribe(reqs, counts, "sub", sent) + assertEquals(0, reqs.activeSubscriptionCount()) + + // The relay answers the REQ after we already sent the CLOSE. + reqs.onIncomingMessage(relay, EoseMessage("sub")) + reqs.onIncomingMessage(relay, ClosedMessage("sub", "closed")) + counts.onIncomingMessage(relay, ClosedMessage("sub", "closed")) + + assertEquals(0, reqs.activeSubscriptionCount(), "late frames must not recreate a dead sub") + assertEquals(0, counts.activeQueryCount()) + assertEquals(0, callbacks, "the listener is gone; late frames must not fire it") + } + + /** + * The registry is the *believed relay state*, deliberately distinct from the desired + * state: while a sub is still desired its row survives everything — including a + * disconnect, which wipes the per-connection wire state but keeps `lastKnownFilters` + * so events still arriving can be linked to the filters the relay was actually + * running. The leak fix must not touch that, and doesn't: it only releases rows + * whose sub is no longer desired. + */ + @Test + fun aStillDesiredSubKeepsLinkingLateEventsToTheFiltersTheRelayWasRunning() = + runTest { + val reqs = PoolRequests() + val counts = PoolCounts() + val sent = mutableListOf() + val relay = FakeRelayClient(relayA) + val subFilters = filters() + + val seen = mutableListOf?>() + val listener = + object : SubscriptionListener { + override suspend fun onEvent( + event: Event, + isLive: Boolean, + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + seen.add(forFilters) + } + } + + val add = reqs.addOrUpdate("sub", mapOf(relayA to subFilters), listener) + reqs.sendToRelayIfChanged("sub", add) { r, cmd -> + sent.add(cmd) + reqs.onSent(r, cmd) + } + + // The socket drops: per-connection wire state is wiped, the row is not. + reqs.onDisconnected(relayA) + + // An event that was already on the wire still resolves to the filters the + // relay was running for this sub. + reqs.onIncomingMessage(relay, EventMessage("sub", event())) + + assertEquals(1, reqs.activeSubscriptionCount(), "a desired sub keeps its row across a disconnect") + assertEquals(1, seen.size, "the listener is still wired") + assertTrue(seen[0] === subFilters, "the late event is linked to the filters the relay was running") + } + + /** + * The counterpart: once a sub is *removed*, `remove()` drops its listener in the + * same call, so nothing is left to process a late frame with. Holding the row past + * that point buys no behaviour — it only inflates the map every connect scans. This + * passes both before and after the fix; it is here to keep that true. + */ + @Test + fun aRemovedSubHasNoListenerLeftToServeLateFrames() = + runTest { + val reqs = PoolRequests() + val counts = PoolCounts() + val sent = mutableListOf() + val relay = FakeRelayClient(relayA) + + var callbacks = 0 + val listener = + object : SubscriptionListener { + override suspend fun onEvent( + event: Event, + isLive: Boolean, + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + callbacks++ + } + + override fun onEose( + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + callbacks++ + } + } + + subscribe(reqs, counts, "sub", setOf(relayA), sent, listener) + unsubscribe(reqs, counts, "sub", sent) + + reqs.onIncomingMessage(relay, EventMessage("sub", event())) + reqs.onIncomingMessage(relay, EoseMessage("sub")) + + assertEquals(0, callbacks, "remove() already dropped the listener; the row has no consumer") + } + + /** The scan cost of a relay lifecycle event must track live subs, not history. */ + @Test + fun connectScanStaysBoundedByLiveSubscriptions() { + val reqs = PoolRequests() + val counts = PoolCounts() + val sent = mutableListOf() + + subscribe(reqs, counts, "tail", setOf(relayA), sent) + + repeat(5000) { i -> + val subId = "churn-$i" + subscribe(reqs, counts, subId, setOf(relayA), sent) + unsubscribe(reqs, counts, subId, sent) + } + + assertEquals(1, reqs.activeSubscriptionCount(), "only the long-lived tail should remain") + + reqs.onConnecting(relayA) + reqs.onDisconnected(relayA) + reqs.onCannotConnect(relayA, "nope") + + assertTrue(reqs.getSubscriptionFiltersOrNull("tail") != null, "the live tail must survive the churn") + } +} From 18bd078a1093297a8104801013b1806c7cd07e25 Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 19 Aug 2026 22:30:29 +0000 Subject: [PATCH 2/2] fix(quartz): release the sub row on unsubscribe even while the client is inactive MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Audit of every subscription teardown path in the app found the previous commit's release was reachable only while the client was ACTIVE: NostrClient.unsubscribe gated the whole flush on isActive(), and that flush is what releases the row. The Android app backgrounds constantly (disconnect() -> isActive() == false) and tears feeds down while it is down there, so every teardown in the background leaked a row for the whole stretch, and all of them were then scanned on every relay connect once the app came back. The flush now always runs; only the wire traffic stays gated. While inactive the pool is disconnected, so sending is both impossible and pointless (the relay dropped the subscription when the socket died), and routing through the pool would materialize a relay client for a relay we deliberately stopped talking to. Also hardens the release against a subscribe of the same id racing an unsubscribe: sendToRelayIfChanged now (re)creates the row when the sub is desired, and re-creates it after a removal that a concurrent subscribe undercut. An empty row means "nothing in flight", so the REQ is still decided and sent — without this, the loser of that race could leave a wanted feed permanently silent, which is far worse than a leaked row. Adds NostrClientSubscriptionLifecycleTest: ten end-to-end cases driving a real NostrClient through a recording socket, each asserting BOTH the wire traffic and the registry size, because it is easy to fix the leak by dropping a frame. They cover subscribe/unsubscribe, teardown while inactive, background churn, re-subscribing the same id (what SubscriptionController does when a feed's relay set empties and refills), dismissal through SubscriptionController, filter changes on a live sub, replay after reconnect, COUNT queries, multi-relay CLOSEs, and survival of the relay-wide refusal block across re-subscribe cycles. Nine of the ten fail on the pre-fix code; the tenth (a filter change on a live sub) passes both ways and is there to catch the opposite regression. Verified side by side that fixed and unfixed produce identical frames on the wire — REQ, REQ, CLOSE, then only the live sub replayed on reconnect — with the registry going from 2 requests + 1 count to 1 + 0. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01WHsz9fXDYysZcKZXhff4XY --- .../nip01Core/relay/client/NostrClient.kt | 19 +- .../nip01Core/relay/client/pool/PoolCounts.kt | 3 + .../relay/client/pool/PoolRequests.kt | 19 +- .../NostrClientSubscriptionLifecycleTest.kt | 454 ++++++++++++++++++ .../pool/PoolRequestsRegistryLeakTest.kt | 60 +++ 5 files changed, 547 insertions(+), 8 deletions(-) create mode 100644 quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClientSubscriptionLifecycleTest.kt diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt index 95c5cda287..5b75850ea3 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClient.kt @@ -336,10 +336,21 @@ class NostrClient( val relaysToUpdateReqs = activeRequests.remove(subId) val relaysToUpdateCounts = activeCounts.remove(subId) - if (isActive()) { - activeRequests.sendToRelayIfChanged(subId, relaysToUpdateReqs, relayPool::sendIfConnected) - activeCounts.sendToRelayIfChanged(subId, relaysToUpdateCounts, relayPool::sendIfConnected) - } + // The flush must run even while INACTIVE, because it is also what releases the + // subscription's state row. The app backgrounds constantly and tears feeds down + // while it is down there; gating this on isActive() would leak a row per teardown + // for the whole background stretch, and every one of them would then be scanned + // on every relay connect once the app comes back. + // + // What is gated is the wire traffic: while inactive the pool is disconnected, so + // sending is both impossible and pointless (the relay dropped the subscription + // when the socket died) -- and routing through the pool would materialize a relay + // client for a relay we deliberately stopped talking to. + val send: (NormalizedRelayUrl, Command) -> Unit = + if (isActive()) relayPool::sendIfConnected else { _, _ -> } + + activeRequests.sendToRelayIfChanged(subId, relaysToUpdateReqs, send) + activeCounts.sendToRelayIfChanged(subId, relaysToUpdateCounts, send) } override fun syncFilters(relay: IRelayClient) { diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolCounts.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolCounts.kt index 34aefc19cd..34bca84f3f 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolCounts.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolCounts.kt @@ -194,6 +194,9 @@ class PoolCounts { // Query is gone and every relay has been told: drop the row (see [relayState]). if (!queries.containsKey(queryId)) { relayState.remove(queryId) + // A concurrent count() may have re-added this id; leave it the empty row it + // expects rather than a missing one. + if (queries.containsKey(queryId)) subState(queryId) } } 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 11d8c38e77..4630a5c58b 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 @@ -423,10 +423,17 @@ class PoolRequests( relaysToUpdate: Set, sync: (NormalizedRelayUrl, Command) -> Unit, ) { - // Never create a row here. NostrClient.unsubscribe calls this on BOTH registries - // with the same id, so a COUNT id would otherwise materialize a permanent - // (and never-scanned-past) row in the REQ registry, and vice-versa. - val state = relayState.get(subId) ?: return + // Don't create a row for an id that isn't ours: NostrClient.unsubscribe calls this + // on BOTH registries with the same id, so a COUNT id would otherwise materialize a + // permanent (and never-scanned-past) row in the REQ registry, and vice-versa. + // + // A DESIRED sub always gets one, though: [addOrUpdate] created it, and if a + // concurrent unsubscribe of the same id released it in between, re-creating is + // exactly right -- an empty row means "nothing in flight", so the REQ below is + // still decided and sent. Returning early there would leave a wanted feed silent. + val state = + relayState.get(subId) + ?: if (desiredSubs.containsKey(subId)) subState(subId) else return relaysToUpdate.forEach { relay -> // Decide + pre-mark atomically under the sub's lock, then send outside it. @@ -442,6 +449,10 @@ class PoolRequests( // and no listener, and correctly no-op. if (!desiredSubs.containsKey(subId)) { relayState.remove(subId) + // A concurrent subscribe() may have re-added this id between the check and + // the removal. Leave it the empty row it expects rather than a missing one: + // empty means "nothing in flight", so its REQ is still decided and sent. + if (desiredSubs.containsKey(subId)) subState(subId) } } diff --git a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClientSubscriptionLifecycleTest.kt b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClientSubscriptionLifecycleTest.kt new file mode 100644 index 0000000000..112f1b2437 --- /dev/null +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/NostrClientSubscriptionLifecycleTest.kt @@ -0,0 +1,454 @@ +/* + * 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.client + +import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener +import com.vitorpamplona.quartz.nip01Core.relay.client.subscriptions.SubscriptionController +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocket +import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocketListener +import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebsocketBuilder +import kotlinx.coroutines.ExperimentalCoroutinesApi +import kotlinx.coroutines.test.TestScope +import kotlinx.coroutines.test.advanceTimeBy +import kotlinx.coroutines.test.runCurrent +import kotlinx.coroutines.test.runTest +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertTrue + +/** + * End-to-end audit of the subscription lifecycle through a real [NostrClient] and a + * recording socket, added alongside the registry-release fix in [PoolRequests]. + * + * Two things must hold together, and it is easy to fix one by breaking the other: + * - **nothing leaks**: the pool's state registries must return to zero once the app + * stops wanting a subscription, on every path that can drop one; + * - **nothing is dropped**: every frame the relay needs to stay in sync — the REQ, the + * CLOSE, the replay after a reconnect — must still reach the wire. + * + * So every test here asserts BOTH the wire traffic and the registry size. + */ +@OptIn(ExperimentalCoroutinesApi::class) +class NostrClientSubscriptionLifecycleTest { + private val url = NormalizedRelayUrl("wss://lifecycle.example.com") + private val url2 = NormalizedRelayUrl("wss://lifecycle2.example.com") + + private class RecordingSocket( + val sent: MutableList, + ) : WebSocket { + override fun needsReconnect() = false + + override fun connect() {} + + override fun disconnect() {} + + override fun send(msg: String): Boolean { + sent.add(msg) + return true + } + } + + /** + * Tracks every relay separately: the pool drops a relay as soon as nothing wants it + * and redials it when something does, so a single `lastListener` would only ever + * describe the most recent socket and quietly hide frames sent to the others. + */ + private class RecordingBuilder : WebsocketBuilder { + val sent = mutableListOf() + val perRelay = mutableMapOf>() + val listeners = mutableMapOf() + + /** Sockets dialed but not yet driven to onOpen — drained by `openAll`. */ + val unopened = mutableListOf() + var connectAttempts = 0 + + override fun build( + url: NormalizedRelayUrl, + out: WebSocketListener, + ): WebSocket { + connectAttempts++ + listeners[url] = out + unopened.add(out) + val perRelaySink = perRelay.getOrPut(url) { mutableListOf() } + return RecordingSocket( + object : MutableList by sent { + override fun add(element: String): Boolean { + perRelaySink.add(element) + return sent.add(element) + } + }, + ) + } + } + + private fun List.reqs() = count { it.startsWith("[\"REQ\"") } + + private fun List.closes() = count { it.startsWith("[\"CLOSE\"") } + + private fun List.counts() = count { it.startsWith("[\"COUNT\"") } + + /** + * Brings every socket the pool has dialed but not yet opened to a ready state. + * Only newly dialed sockets are opened: re-opening a live one would make the client + * treat it as a fresh connection and replay every filter, inflating REQ counts. + */ + private fun TestScope.openAll(builder: RecordingBuilder) { + // Loop: opening one relay can make the pool dial another. + repeat(3) { + val batch = builder.unopened.toList() + builder.unopened.clear() + if (batch.isEmpty() && it > 0) return + batch.forEach { listener -> listener.onOpen(50, false) } + advanceTimeBy(500) + runCurrent() + } + } + + private fun TestScope.settle() { + advanceTimeBy(1_000) + runCurrent() + } + + private fun filters() = listOf(Filter(kinds = listOf(1))) + + @Test + fun subscribeThenUnsubscribeSendsBothFramesAndReleasesTheRegistry() = + runTest { + val builder = RecordingBuilder() + val client = NostrClient(builder, this) + try { + advanceTimeBy(500) + runCurrent() + + client.subscribe("sub", mapOf(url to filters())) + openAll(builder) + assertEquals(1, builder.sent.reqs(), "the REQ must reach the relay") + assertEquals(1, client.registrySizes().liveRequests, "a live sub is tracked") + + client.unsubscribe("sub") + advanceTimeBy(500) + runCurrent() + + assertEquals(1, builder.sent.closes(), "the relay must be told to CLOSE") + assertEquals(0, client.registrySizes().liveRequests, "the row must be released") + assertEquals(0, client.registrySizes().liveCounts) + } finally { + client.close() + } + } + + /** + * The app backgrounds (host calls [NostrClient.disconnect]) and a ViewModel or a + * composable then tears its subscription down. No CLOSE can be sent — the socket is + * gone and the relay already dropped the sub — but the row must still be released, + * or it survives the whole background stretch and inflates every later connect scan. + */ + @Test + fun unsubscribeWhileInactiveStillReleasesTheRegistry() = + runTest { + val builder = RecordingBuilder() + val client = NostrClient(builder, this) + try { + advanceTimeBy(500) + runCurrent() + + client.subscribe("sub", mapOf(url to filters())) + openAll(builder) + assertEquals(1, client.registrySizes().liveRequests) + + client.disconnect() + advanceTimeBy(500) + runCurrent() + + client.unsubscribe("sub") + advanceTimeBy(500) + runCurrent() + + assertEquals(0, client.registrySizes().liveRequests, "an inactive client must still release the row") + } finally { + client.close() + } + } + + /** Many subs torn down while backgrounded must not accumulate. */ + @Test + fun churnWhileInactiveDoesNotAccumulate() = + runTest { + val builder = RecordingBuilder() + val client = NostrClient(builder, this) + try { + advanceTimeBy(500) + runCurrent() + client.disconnect() + advanceTimeBy(500) + runCurrent() + + repeat(200) { i -> + client.subscribe("bg-$i", mapOf(url to filters())) + client.unsubscribe("bg-$i") + } + advanceTimeBy(500) + runCurrent() + + assertEquals(0, client.registrySizes().liveRequests, "background churn must not accumulate rows") + } finally { + client.close() + } + } + + /** + * The core app pattern: [SubscriptionController.updateRelaysIfNeeded] unsubscribes a + * feed whose relay set went empty and re-subscribes THE SAME id when it comes back. + * The second subscribe must produce a fresh REQ — a stale row that made the client + * think a REQ was already in flight would silently leave the feed empty. + */ + @Test + fun resubscribingTheSameIdAfterUnsubscribeSendsAFreshReq() = + runTest { + val builder = RecordingBuilder() + val client = NostrClient(builder, this) + try { + advanceTimeBy(500) + runCurrent() + + repeat(5) { + client.subscribe("feed", mapOf(url to filters())) + // The pool drops the relay when nothing wants it and redials on the + // next subscribe, so each cycle needs its socket brought up. + openAll(builder) + client.unsubscribe("feed") + settle() + } + + assertEquals(5, builder.sent.reqs(), "every re-subscribe must re-REQ") + assertEquals(5, builder.sent.closes(), "every unsubscribe must CLOSE") + assertEquals(0, client.registrySizes().liveRequests) + } finally { + client.close() + } + } + + /** Same cycle driven through SubscriptionController, which is what the app uses. */ + @Test + fun subscriptionControllerDismissAndRecreateStaysInSync() = + runTest { + val builder = RecordingBuilder() + val client = NostrClient(builder, this) + try { + val controller = SubscriptionController(client) + advanceTimeBy(500) + runCurrent() + + repeat(4) { + val sub = controller.requestNewSubscription("feed", object : SubscriptionListener {}) + client.subscribe("feed", mapOf(url to filters()), sub.listener) + openAll(builder) + controller.dismissSubscription(sub) + settle() + } + + assertEquals(4, builder.sent.reqs()) + assertEquals(4, builder.sent.closes()) + assertEquals(0, client.registrySizes().liveRequests, "dismissed subs must not linger") + } finally { + client.close() + } + } + + /** A live sub whose filters change keeps exactly one row and re-REQs. */ + @Test + fun filterChangesOnALiveSubKeepOneRowAndResend() = + runTest { + val builder = RecordingBuilder() + val client = NostrClient(builder, this) + try { + advanceTimeBy(500) + runCurrent() + + client.subscribe("feed", mapOf(url to listOf(Filter(kinds = listOf(1))))) + openAll(builder) + builder.listeners.getValue(url).onMessage("[\"EOSE\",\"feed\"]") + advanceTimeBy(500) + runCurrent() + + client.subscribe("feed", mapOf(url to listOf(Filter(kinds = listOf(7))))) + advanceTimeBy(500) + runCurrent() + + assertEquals(1, client.registrySizes().liveRequests, "a filter change is one sub, not two") + assertTrue(builder.sent.reqs() >= 2, "the changed filters must be re-REQed") + assertEquals(0, builder.sent.closes(), "a filter change must not CLOSE the sub") + } finally { + client.close() + } + } + + /** A reconnect must replay every live sub — and only the live ones. */ + @Test + fun reconnectReplaysLiveSubsOnly() = + runTest { + val builder = RecordingBuilder() + val client = NostrClient(builder, this) + try { + advanceTimeBy(500) + runCurrent() + + client.subscribe("keep", mapOf(url to filters())) + client.subscribe("drop", mapOf(url to filters())) + openAll(builder) + advanceTimeBy(500) + runCurrent() + + client.unsubscribe("drop") + advanceTimeBy(500) + runCurrent() + assertEquals(1, client.registrySizes().liveRequests) + + builder.sent.clear() + builder.listeners.getValue(url).onClosed(1000, "server closed") + advanceTimeBy(500) + runCurrent() + client.connect() + advanceTimeBy(1_000) + runCurrent() + openAll(builder) + + assertEquals(1, builder.sent.reqs(), "exactly the live sub is replayed") + assertTrue(builder.sent.any { it.contains("\"keep\"") }, "the live sub must be replayed") + assertTrue(builder.sent.none { it.contains("\"drop\"") }, "the dropped sub must not be replayed") + } finally { + client.close() + } + } + + /** A COUNT query must release its row too, and not touch the REQ registry. */ + @Test + fun countQueriesReleaseTheirOwnRegistryOnly() = + runTest { + val builder = RecordingBuilder() + val client = NostrClient(builder, this) + try { + advanceTimeBy(500) + runCurrent() + + client.subscribe("req", mapOf(url to filters())) + client.count("cnt", mapOf(url to filters())) + openAll(builder) + advanceTimeBy(500) + runCurrent() + + assertEquals(1, client.registrySizes().liveRequests) + assertEquals(1, client.registrySizes().liveCounts) + assertEquals(1, builder.sent.counts(), "the COUNT must reach the relay") + + client.unsubscribe("cnt") + advanceTimeBy(500) + runCurrent() + + assertEquals(1, client.registrySizes().liveRequests, "closing a COUNT must not disturb the REQ") + assertEquals(0, client.registrySizes().liveCounts) + + client.unsubscribe("req") + advanceTimeBy(500) + runCurrent() + assertEquals(0, client.registrySizes().liveRequests) + assertEquals(0, client.registrySizes().liveCounts) + } finally { + client.close() + } + } + + /** + * Releasing a sub's row also drops that sub's per-filter refusal memory, so this pins + * what must NOT be dropped: the relay-WIDE capability block lives in + * [com.vitorpamplona.quartz.nip01Core.relay.client.pool.RelayReqRefusals], keyed by + * relay and never by sub, so a relay that refuses to serve reads at all stays blocked + * across any number of unsubscribe/re-subscribe cycles. That is the guard that stops + * a relay being hammered; the per-sub counter only sharpens it. + */ + @Test + fun relayWideRefusalSurvivesUnsubscribeAndResubscribe() = + runTest { + val builder = RecordingBuilder() + val client = NostrClient(builder, this) + try { + settle() + + // Two refusals of the "no reads" class block the relay pool-wide. + repeat(2) { round -> + client.subscribe("feed-$round", mapOf(url to filters())) + openAll(builder) + builder.listeners + .getValue(url) + .onMessage("[\"CLOSED\",\"feed-$round\",\"unsupported: this relay does not accept REQ\"]") + settle() + client.unsubscribe("feed-$round") + settle() + } + + assertEquals(0, client.registrySizes().liveRequests, "cycled subs must not linger") + + // A brand-new sub on the same relay must now be suppressed, even though + // every earlier sub's own row is long gone. + builder.sent.clear() + client.subscribe("feed-new", mapOf(url to filters())) + openAll(builder) + settle() + + assertEquals( + 0, + builder.sent.reqs(), + "a relay blocked pool-wide must not be re-offered a REQ after the sub rows were released", + ) + } finally { + client.close() + } + } + + /** A multi-relay sub must CLOSE on every relay that got a REQ before releasing. */ + @Test + fun multiRelaySubClosesOnEveryRelayBeforeRelease() = + runTest { + val builder = RecordingBuilder() + val client = NostrClient(builder, this) + try { + settle() + + client.subscribe("multi", mapOf(url to filters(), url2 to filters())) + openAll(builder) + settle() + + val reqRelays = builder.perRelay.filterValues { it.reqs() > 0 }.keys + assertEquals(setOf(url, url2), reqRelays, "both relays must get the REQ") + + client.unsubscribe("multi") + settle() + + val closeRelays = builder.perRelay.filterValues { it.closes() > 0 }.keys + assertEquals(reqRelays, closeRelays, "every relay that got a REQ must get a CLOSE") + assertEquals(0, client.registrySizes().liveRequests) + } finally { + client.close() + } + } +} diff --git a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolRequestsRegistryLeakTest.kt b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolRequestsRegistryLeakTest.kt index a3c7ed9c12..76842aea4d 100644 --- a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolRequestsRegistryLeakTest.kt +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolRequestsRegistryLeakTest.kt @@ -342,4 +342,64 @@ class PoolRequestsRegistryLeakTest { assertTrue(reqs.getSubscriptionFiltersOrNull("tail") != null, "the live tail must survive the churn") } + + /** + * A subscribe and an unsubscribe of the SAME id racing on different threads. The + * release must never swallow the new subscription's REQ: a feed the app still wants + * that never gets a REQ is silent forever, which is far worse than a leaked row. + * + * Honest caveat: the window this protects (a release landing between a concurrent + * subscribe's row creation and its send decision) is a few instructions wide, and + * this test does NOT reliably reproduce it — it passes with the guard removed. It is + * a regression guard on the invariant "registry membership tracks desired-ness", not + * proof that the guard fires. The guard is kept because it is two comparisons and the + * failure it prevents is a permanently silent subscription. + */ + @Test + fun subscribeRacingAnUnsubscribeOfTheSameIdStillSendsItsReq() { + repeat(300) { episode -> + val reqs = PoolRequests() + val sent = java.util.concurrent.ConcurrentLinkedQueue() + val subId = "raced-$episode" + + // A live sub about to be torn down. + reqs.addOrUpdate(subId, mapOf(relayA to filters()), null) + reqs.sendToRelayIfChanged(subId, setOf(relayA)) { r, cmd -> reqs.onSent(r, cmd) } + + val start = java.util.concurrent.CountDownLatch(1) + val unsub = + Thread { + start.await() + val relays = reqs.remove(subId) + reqs.sendToRelayIfChanged(subId, relays) { r, cmd -> + sent.add(cmd) + reqs.onSent(r, cmd) + } + } + val resub = + Thread { + start.await() + val relays = reqs.addOrUpdate(subId, mapOf(relayA to filters()), null) + reqs.sendToRelayIfChanged(subId, relays) { r, cmd -> + sent.add(cmd) + reqs.onSent(r, cmd) + } + } + unsub.start() + resub.start() + start.countDown() + unsub.join() + resub.join() + + // Whoever won, the invariant is: if the sub is still desired it must have a + // row (so later filter changes and reconnect replays work), and if it is not + // desired it must have none. + val desired = reqs.getSubscriptionFiltersOrNull(subId) != null + assertEquals( + if (desired) 1 else 0, + reqs.activeSubscriptionCount(), + "episode $episode: registry must match desired-ness (desired=$desired)", + ) + } + } }