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..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) { @@ -458,6 +469,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..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 @@ -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,14 @@ 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) + } } 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..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 @@ -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,18 @@ class PoolRequests( relaysToUpdate: Set, sync: (NormalizedRelayUrl, Command) -> Unit, ) { - val state = subState(subId) + // 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. val cmd = state.withLock(relay) { decideCommandLocked(state, subId, relay) } @@ -406,6 +442,18 @@ 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) + // 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 new file mode 100644 index 0000000000..76842aea4d --- /dev/null +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/pool/PoolRequestsRegistryLeakTest.kt @@ -0,0 +1,405 @@ +/* + * 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") + } + + /** + * 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)", + ) + } + } +}