diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/AppModules.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/AppModules.kt index 393e9b64d8..5cbb2cc366 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/AppModules.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/AppModules.kt @@ -122,6 +122,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.RelayLogger import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.RelayOfflineTracker +import com.vitorpamplona.quartz.nip01Core.relay.client.limits.RelayLimits import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.stats.RelayReqStats import com.vitorpamplona.quartz.nip01Core.relay.client.stats.RelayStats import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.CachingEventDecoder @@ -710,6 +711,9 @@ class AppModules( // Captures statistics about relays val relayStats = RelayStats(client) + // Caches the latest LIMITS (rights + limits) each relay advertises. + val relayLimits = RelayLimits(client) + // Resource-usage ledger: relay traffic/reconnect + connection-time, // foreground-time, process-CPU, and signature-verification collectors. init { diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/limits/RelayLimits.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/limits/RelayLimits.kt new file mode 100644 index 0000000000..36acfaf475 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/limits/RelayLimits.kt @@ -0,0 +1,100 @@ +/* + * 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.limits + +import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient +import com.vitorpamplona.quartz.nip01Core.relay.client.listeners.RelayConnectionListener +import com.vitorpamplona.quartz.nip01Core.relay.client.single.IRelayClient +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.LimitsMessage +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import com.vitorpamplona.quartz.utils.Log +import kotlinx.collections.immutable.PersistentMap +import kotlinx.collections.immutable.persistentMapOf +import kotlinx.coroutines.flow.MutableStateFlow +import kotlinx.coroutines.flow.StateFlow +import kotlinx.coroutines.flow.asStateFlow +import kotlinx.coroutines.flow.update + +/** + * Caches the latest `LIMITS` (relay rights + limits) advertised by each relay. + * + * A relay sends `LIMITS` on connect and again whenever the connection's rights + * change (e.g. after a successful NIP-42 AUTH flips `can_write` on). Like + * [com.vitorpamplona.quartz.nip01Core.relay.client.auth.RelayAuthenticator], + * this is an accessory that registers a [RelayConnectionListener], keeps + * per-relay state, and publishes it as a Compose-stable [StateFlow] — but it is + * purely passive: it never replies to the relay. + * + * The cache is connection-scoped: a relay's entry is dropped on disconnect, so a + * stale limit from a previous session never leaks into a new one. A fresh + * connection re-advertises `LIMITS` on connect. + */ +class RelayLimits( + val client: INostrClient, +) { + // onIncomingMessage / onDisconnected fire on the per-relay socket dispatcher + // thread, so this is mutated concurrently. MutableStateFlow.update is an + // atomic compare-and-set loop over an immutable PersistentMap, so concurrent + // writers from different relays never corrupt the map or lose an update. + private val _limitsFlow = MutableStateFlow>(persistentMapOf()) + + /** + * Per-relay `LIMITS` as an immutable, Compose-stable snapshot map. The map + * identity changes on every mutation, so downstream + * [kotlinx.coroutines.flow.distinctUntilChanged] and Compose `@Immutable` + * skipping both work correctly. + */ + val limitsFlow: StateFlow> = _limitsFlow.asStateFlow() + + /** The most recent `LIMITS` the relay advertised, or null if none seen on the current connection. */ + fun get(url: NormalizedRelayUrl): LimitsMessage? = _limitsFlow.value[url] + + fun snapshot(): Map = _limitsFlow.value + + private val clientListener = + object : RelayConnectionListener { + override fun onIncomingMessage( + relay: IRelayClient, + msgStr: String, + msg: Message, + ) { + if (msg is LimitsMessage) { + _limitsFlow.update { it.putting(relay.url, msg) } + } + } + + override fun onDisconnected(relay: IRelayClient) { + _limitsFlow.update { it.removing(relay.url) } + } + } + + init { + Log.d("RelayLimits", "Init, Subscribe") + client.addConnectionListener(clientListener) + } + + fun destroy() { + // makes sure to run + Log.d("RelayLimits", "Destroy, Unsubscribe") + client.removeConnectionListener(clientListener) + } +} diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/commands/toClient/LimitsMessage.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/commands/toClient/LimitsMessage.kt index cb2f68f2f6..38b692825b 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/commands/toClient/LimitsMessage.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/commands/toClient/LimitsMessage.kt @@ -20,8 +20,10 @@ */ package com.vitorpamplona.quartz.nip01Core.relay.commands.toClient +import androidx.compose.runtime.Immutable + /** - * NIP-22 `LIMITS` message: the relay advertises the current rights and limits + * `LIMITS` message: the relay advertises the current rights and limits * for this connection. Sent upon connection and again at any point during the * connection to reflect the client's changing rights (e.g. after a NIP-42 AUTH). * @@ -32,6 +34,7 @@ package com.vitorpamplona.quartz.nip01Core.relay.commands.toClient * any previously cached value and not assume a default. Clients cache this * payload on the relay connection and apply it when sending `EVENT`s and `REQ`s. */ +@Immutable class LimitsMessage( /** Whether clients may publish events to this relay. */ val canWrite: Boolean? = null, diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/limits/RelayLimitsTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/limits/RelayLimitsTest.kt new file mode 100644 index 0000000000..b88f2e2dfc --- /dev/null +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/limits/RelayLimitsTest.kt @@ -0,0 +1,136 @@ +/* + * 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.limits + +import com.vitorpamplona.quartz.nip01Core.relay.client.EmptyNostrClient +import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient +import com.vitorpamplona.quartz.nip01Core.relay.client.listeners.RelayConnectionListener +import com.vitorpamplona.quartz.nip01Core.relay.client.single.IRelayClient +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.EoseMessage +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.LimitsMessage +import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.Command +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertNull +import kotlin.test.assertTrue + +class RelayLimitsTest { + private class CapturingClient( + private val delegate: INostrClient = EmptyNostrClient(), + ) : INostrClient by delegate { + var captured: RelayConnectionListener? = null + + override fun addConnectionListener(listener: RelayConnectionListener) { + captured = listener + } + } + + 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 sendOrConnectAndSync(cmd: Command) = Unit + + override fun sendIfConnected(cmd: Command) = Unit + + override fun disconnect() = Unit + } + + private fun setup(): Pair { + val client = CapturingClient() + val limits = RelayLimits(client) + val listener = client.captured ?: error("RelayLimits did not register a listener") + return limits to listener + } + + @Test + fun cachesLimitsPerRelay() { + val (limits, listener) = setup() + val relay = FakeRelayClient(NormalizedRelayUrl("wss://relay.example/")) + + assertNull(limits.get(relay.url), "No limits before any LIMITS message") + + listener.onIncomingMessage(relay, "", LimitsMessage(canWrite = true, maxLimit = 200)) + + assertEquals(true, limits.get(relay.url)?.canWrite) + assertEquals(200, limits.get(relay.url)?.maxLimit) + assertEquals(limits.get(relay.url), limits.limitsFlow.value[relay.url]) + } + + @Test + fun laterLimitsReplaceEarlierOnes() { + val (limits, listener) = setup() + val relay = FakeRelayClient(NormalizedRelayUrl("wss://relay.example/")) + + listener.onIncomingMessage(relay, "", LimitsMessage(canWrite = false, maxLimit = 200)) + listener.onIncomingMessage(relay, "", LimitsMessage(canWrite = true, maxLimit = 500)) + + // A relay re-advertises LIMITS when rights change (e.g. after AUTH flips can_write). + assertEquals(true, limits.get(relay.url)?.canWrite) + assertEquals(500, limits.get(relay.url)?.maxLimit) + } + + @Test + fun tracksLimitsForDistinctRelaysIndependently() { + val (limits, listener) = setup() + val relayA = FakeRelayClient(NormalizedRelayUrl("wss://a.example/")) + val relayB = FakeRelayClient(NormalizedRelayUrl("wss://b.example/")) + + listener.onIncomingMessage(relayA, "", LimitsMessage(maxLimit = 100)) + listener.onIncomingMessage(relayB, "", LimitsMessage(maxLimit = 999)) + + assertEquals(100, limits.get(relayA.url)?.maxLimit) + assertEquals(999, limits.get(relayB.url)?.maxLimit) + assertEquals(2, limits.snapshot().size) + } + + @Test + fun dropsCachedLimitsOnDisconnect() { + val (limits, listener) = setup() + val relay = FakeRelayClient(NormalizedRelayUrl("wss://relay.example/")) + + listener.onIncomingMessage(relay, "", LimitsMessage(canRead = true)) + assertTrue(limits.get(relay.url) != null) + + listener.onDisconnected(relay) + assertNull(limits.get(relay.url), "Limits are connection-scoped and cleared on disconnect") + assertTrue(limits.snapshot().isEmpty()) + } + + @Test + fun ignoresNonLimitsMessages() { + val (limits, listener) = setup() + val relay = FakeRelayClient(NormalizedRelayUrl("wss://relay.example/")) + + listener.onIncomingMessage(relay, "", EoseMessage("sub1")) + + assertNull(limits.get(relay.url)) + assertTrue(limits.snapshot().isEmpty()) + } +}