mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 03:38:23 +00:00
feat: cache and expose relay LIMITS per relay
Adds RelayLimits, a passive connection-listener accessory (modeled on RelayAuthenticator) that caches the latest LIMITS each relay advertises and publishes it as a Compose-stable StateFlow, so consumers can read a relay's current rights/limits instead of only observing the raw message. - RelayLimits: caches LimitsMessage per NormalizedRelayUrl, exposes limitsFlow (StateFlow), get(url) and snapshot(); connection-scoped (entry dropped on disconnect so stale limits don't leak). - LimitsMessage marked @Immutable for Compose stability in the flow map. - Wired into AppModules next to relayStats (Amethyst.instance.relayLimits). - Tests: per-relay caching, later-replaces-earlier, independent relays, drop-on-disconnect, and ignore non-LIMITS messages. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01464jkunWPtYhTReoc3fUQQ
This commit is contained in:
@@ -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 {
|
||||
|
||||
+100
@@ -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<PersistentMap<NormalizedRelayUrl, LimitsMessage>>(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<PersistentMap<NormalizedRelayUrl, LimitsMessage>> = _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<NormalizedRelayUrl, LimitsMessage> = _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)
|
||||
}
|
||||
}
|
||||
+4
-1
@@ -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,
|
||||
|
||||
+136
@@ -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<RelayLimits, RelayConnectionListener> {
|
||||
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())
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user