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)", + ) + } + } }