fix(quartz): release the sub row on unsubscribe even while the client is inactive

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 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01WHsz9fXDYysZcKZXhff4XY
This commit is contained in:
Claude
2026-08-19 22:30:29 +00:00
parent ca89b30fbe
commit 18bd078a10
5 changed files with 547 additions and 8 deletions
@@ -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) {
@@ -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)
}
}
@@ -423,10 +423,17 @@ class PoolRequests(
relaysToUpdate: Set<NormalizedRelayUrl>,
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)
}
}
@@ -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<String>,
) : 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<String>()
val perRelay = mutableMapOf<NormalizedRelayUrl, MutableList<String>>()
val listeners = mutableMapOf<NormalizedRelayUrl, WebSocketListener>()
/** Sockets dialed but not yet driven to onOpen — drained by `openAll`. */
val unopened = mutableListOf<WebSocketListener>()
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<String> by sent {
override fun add(element: String): Boolean {
perRelaySink.add(element)
return sent.add(element)
}
},
)
}
}
private fun List<String>.reqs() = count { it.startsWith("[\"REQ\"") }
private fun List<String>.closes() = count { it.startsWith("[\"CLOSE\"") }
private fun List<String>.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()
}
}
}
@@ -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<Command>()
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)",
)
}
}
}