mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 11:48:24 +00:00
fix(quartz): release pool request/count state rows when a sub is removed
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 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01WHsz9fXDYysZcKZXhff4XY
This commit is contained in:
+18
@@ -143,6 +143,18 @@ interface INostrClient : AutoCloseable {
|
||||
|
||||
fun activeCounts(url: NormalizedRelayUrl): Map<String, List<Filter>>
|
||||
|
||||
/**
|
||||
* 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<HexKey>
|
||||
|
||||
/** 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,
|
||||
)
|
||||
|
||||
+6
@@ -458,6 +458,12 @@ class NostrClient(
|
||||
|
||||
override fun activeCounts(url: NormalizedRelayUrl): Map<String, List<Filter>> = activeCounts.activeFiltersFor(url)
|
||||
|
||||
override fun registrySizes() =
|
||||
RegistrySizes(
|
||||
liveRequests = activeRequests.activeSubscriptionCount(),
|
||||
liveCounts = activeCounts.activeQueryCount(),
|
||||
)
|
||||
|
||||
override fun activeOutboxCache(url: NormalizedRelayUrl): Set<HexKey> = eventOutbox.activeOutboxCacheFor(url)
|
||||
|
||||
override fun activeOutboxEvents(url: NormalizedRelayUrl): List<Event> = eventOutbox.activeOutboxEventsFor(url)
|
||||
|
||||
+33
-5
@@ -38,9 +38,26 @@ class PoolCounts {
|
||||
private var queries = LargeCache<String, Map<NormalizedRelayUrl, List<Filter>>>()
|
||||
val relays = MutableStateFlow(setOf<NormalizedRelayUrl>())
|
||||
|
||||
/**
|
||||
* **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<String, CountQueryState<NormalizedRelayUrl>>()
|
||||
|
||||
fun subState(subId: String): CountQueryState<NormalizedRelayUrl> = relayState.getOrCreate(subId) { CountQueryState() }
|
||||
private fun subState(subId: String): CountQueryState<NormalizedRelayUrl> = 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<NormalizedRelayUrl>? = 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<NormalizedRelayUrl>()
|
||||
@@ -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(
|
||||
|
||||
+41
-4
@@ -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<String, RequestSubscriptionState<NormalizedRelayUrl>>()
|
||||
|
||||
fun subState(subId: String): RequestSubscriptionState<NormalizedRelayUrl> = relayState.getOrCreate(subId) { RequestSubscriptionState() }
|
||||
private fun subState(subId: String): RequestSubscriptionState<NormalizedRelayUrl> = 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<NormalizedRelayUrl>,
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
+345
@@ -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<Command>,
|
||||
) {
|
||||
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<NormalizedRelayUrl>,
|
||||
sent: MutableList<Command> = 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<Command>()
|
||||
|
||||
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<Command>()
|
||||
|
||||
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<Command>()
|
||||
|
||||
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<Command>()
|
||||
val relay = FakeRelayClient(relayA)
|
||||
|
||||
var callbacks = 0
|
||||
val listener =
|
||||
object : SubscriptionListener {
|
||||
override fun onEose(
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
callbacks++
|
||||
}
|
||||
|
||||
override fun onClosed(
|
||||
message: String,
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
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<Command>()
|
||||
val relay = FakeRelayClient(relayA)
|
||||
val subFilters = filters()
|
||||
|
||||
val seen = mutableListOf<List<Filter>?>()
|
||||
val listener =
|
||||
object : SubscriptionListener {
|
||||
override suspend fun onEvent(
|
||||
event: Event,
|
||||
isLive: Boolean,
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
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<Command>()
|
||||
val relay = FakeRelayClient(relayA)
|
||||
|
||||
var callbacks = 0
|
||||
val listener =
|
||||
object : SubscriptionListener {
|
||||
override suspend fun onEvent(
|
||||
event: Event,
|
||||
isLive: Boolean,
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
callbacks++
|
||||
}
|
||||
|
||||
override fun onEose(
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
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<Command>()
|
||||
|
||||
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")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user