mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 19:53:08 +00:00
feat(quartz): relay connection observability + stable connection ids
Gives relay operators metrics/logging hooks without patching the engine, and removes the hashCode()-keyed connection registry the audit flagged. - RelaySession gains a stable, process-unique `id` (monotonic counter). - RelayConnectionListener (onConnect/onDisconnect, no-op default) is accepted by NostrServer and ReqResponderServer; both now key their connection registry by `id` instead of hashCode() (no more identity-collision hole). - Both servers expose a live `activeConnections` gauge. Teardown accounting is idempotent (double close counted once) and onDisconnect fires for any connections still open at server close. - Tests + RELAY.md Observability section. https://claude.ai/code/session_016YBS2pWCBSDgAMthHzfCTr
This commit is contained in:
+26
-1
@@ -312,6 +312,30 @@ class MyRelayTest {
|
||||
}
|
||||
```
|
||||
|
||||
## Observability
|
||||
|
||||
Both servers take an optional `RelayConnectionListener` and expose a live
|
||||
`activeConnections` gauge. Connections carry a stable, process-unique
|
||||
`RelaySession.id` (used to key the server's registry and the listener
|
||||
callbacks), so you can correlate the open/close of the same connection in logs
|
||||
and metrics:
|
||||
|
||||
```kotlin
|
||||
val server = NostrServer(
|
||||
store = store,
|
||||
listener = object : RelayConnectionListener {
|
||||
override fun onConnect(connectionId: Long) = metrics.connections.inc()
|
||||
override fun onDisconnect(connectionId: Long) = metrics.connections.dec()
|
||||
},
|
||||
)
|
||||
server.activeConnections // Long, current count
|
||||
```
|
||||
|
||||
`onConnect`/`onDisconnect` fire at most once per connection (a double `close()`
|
||||
is accounted once) and `onDisconnect` is also fired for any connections still
|
||||
open when the server itself is closed. Callbacks can run on any transport
|
||||
coroutine — keep them cheap and non-blocking.
|
||||
|
||||
## Key Source Files
|
||||
|
||||
```
|
||||
@@ -321,7 +345,8 @@ quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/
|
||||
│ ├── ReqResponderServer.kt # Non-storage engine (search/redirector/computed)
|
||||
│ ├── ReqResponder.kt # Flow<Event> REQ-responder SPI
|
||||
│ ├── SessionBackend.kt # Data-plane seam (LiveEventStore / ReqResponderBackend)
|
||||
│ ├── RelaySession.kt # Per-connection handler
|
||||
│ ├── RelayConnectionListener.kt # Connection open/close observability hook
|
||||
│ ├── RelaySession.kt # Per-connection handler (stable .id)
|
||||
│ ├── LiveEventStore.kt # Reactive event streaming (storage SessionBackend)
|
||||
│ ├── IRelayPolicy.kt # Policy interface + PolicyResult + onAuthenticated
|
||||
│ └── policies/
|
||||
|
||||
+42
-16
@@ -28,6 +28,8 @@ import com.vitorpamplona.quartz.utils.cache.LargeCache
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlin.concurrent.atomics.AtomicLong
|
||||
import kotlin.concurrent.atomics.ExperimentalAtomicApi
|
||||
import kotlin.coroutines.CoroutineContext
|
||||
|
||||
/**
|
||||
@@ -47,17 +49,24 @@ import kotlin.coroutines.CoroutineContext
|
||||
* @param negentropySettings NIP-77 server-side tuning (frame cap,
|
||||
* snapshot cap, per-connection session cap). Defaults to strfry-
|
||||
* parity values; see [NegentropySettings].
|
||||
* @param listener Observability hook fired as connections open and close,
|
||||
* keyed by [RelaySession.id]. Defaults to a no-op.
|
||||
*/
|
||||
@OptIn(ExperimentalAtomicApi::class)
|
||||
class NostrServer(
|
||||
private val store: IEventStore,
|
||||
private val policyBuilder: () -> IRelayPolicy = { VerifyPolicy },
|
||||
private val parentContext: CoroutineContext = SupervisorJob(),
|
||||
parallelVerify: Boolean = false,
|
||||
private val negentropySettings: NegentropySettings = NegentropySettings.Default,
|
||||
private val listener: RelayConnectionListener = RelayConnectionListener.None,
|
||||
) : AutoCloseable {
|
||||
/** Scope for all subscriptions. */
|
||||
private val scope = CoroutineScope(parentContext + SupervisorJob())
|
||||
|
||||
/** Live count of registered connections; backs [activeConnections]. */
|
||||
private val activeCount = AtomicLong(0L)
|
||||
|
||||
/**
|
||||
* Group-commit writer shared across every connected session.
|
||||
* Sessions hand off EVENT publishes here instead of awaiting
|
||||
@@ -74,8 +83,11 @@ class NostrServer(
|
||||
|
||||
private val subStore = LiveEventStore(store, ingest)
|
||||
|
||||
/** Active client sessions keyed by an opaque connection id. */
|
||||
private val connections = LargeCache<Int, RelaySession>()
|
||||
/** Active client sessions keyed by [RelaySession.id]. */
|
||||
private val connections = LargeCache<Long, RelaySession>()
|
||||
|
||||
/** Number of connections currently registered with this server. */
|
||||
val activeConnections: Long get() = activeCount.load()
|
||||
|
||||
/**
|
||||
* Registers a new client connection.
|
||||
@@ -83,19 +95,29 @@ class NostrServer(
|
||||
* @param send Callback the server uses to send JSON messages to this client.
|
||||
* Implementations must be safe to call from any coroutine.
|
||||
*/
|
||||
fun connect(send: (String) -> Unit) =
|
||||
RelaySession(
|
||||
policy = policyBuilder(),
|
||||
store = subStore,
|
||||
scope = scope,
|
||||
onSend = send,
|
||||
onClose = { session ->
|
||||
connections.remove(session.hashCode())
|
||||
},
|
||||
negentropySettings = negentropySettings,
|
||||
).also { session ->
|
||||
connections.put(session.hashCode(), session)
|
||||
}
|
||||
fun connect(send: (String) -> Unit): RelaySession {
|
||||
val session =
|
||||
RelaySession(
|
||||
policy = policyBuilder(),
|
||||
store = subStore,
|
||||
scope = scope,
|
||||
onSend = send,
|
||||
onClose = { closed ->
|
||||
// Idempotent: only account for the first teardown of a
|
||||
// given connection so a double close() can't underflow
|
||||
// the gauge or double-fire the listener.
|
||||
if (connections.remove(closed.id) != null) {
|
||||
activeCount.addAndFetch(-1L)
|
||||
listener.onDisconnect(closed.id)
|
||||
}
|
||||
},
|
||||
negentropySettings = negentropySettings,
|
||||
)
|
||||
connections.put(session.id, session)
|
||||
activeCount.addAndFetch(1L)
|
||||
listener.onConnect(session.id)
|
||||
return session
|
||||
}
|
||||
|
||||
/**
|
||||
* Registers a new client connection and serves it for the duration of
|
||||
@@ -121,8 +143,12 @@ class NostrServer(
|
||||
* Shuts down the server, cancelling all subscriptions and closing the store.
|
||||
*/
|
||||
override fun close() {
|
||||
connections.forEach { _, session -> session.cancelAllSubscriptions() }
|
||||
connections.forEach { _, session ->
|
||||
session.cancelAllSubscriptions()
|
||||
listener.onDisconnect(session.id)
|
||||
}
|
||||
connections.clear()
|
||||
activeCount.store(0L)
|
||||
ingest.close()
|
||||
scope.cancel()
|
||||
store.close()
|
||||
|
||||
+50
@@ -0,0 +1,50 @@
|
||||
/*
|
||||
* 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.server
|
||||
|
||||
/**
|
||||
* Observability hook for the relay server connection lifecycle. Both
|
||||
* [NostrServer] and [ReqResponderServer] accept one and invoke it as
|
||||
* connections open and close, keyed by the stable per-connection
|
||||
* [RelaySession.id], so an operator can drive metrics (active-connection
|
||||
* gauges, churn counters) and per-connection logging without patching the
|
||||
* engine.
|
||||
*
|
||||
* Both callbacks default to no-ops; override only what you need. They may fire
|
||||
* from any transport coroutine, so implementations must be thread-safe and
|
||||
* cheap — do not block (no synchronous I/O); hand off to your metrics/logging
|
||||
* pipeline and return.
|
||||
*/
|
||||
interface RelayConnectionListener {
|
||||
/** A new connection was registered. [connectionId] is [RelaySession.id]. */
|
||||
fun onConnect(connectionId: Long) {}
|
||||
|
||||
/**
|
||||
* A connection was torn down (client CLOSE, transport drop, or server
|
||||
* shutdown). Fires at most once per connection.
|
||||
*/
|
||||
fun onDisconnect(connectionId: Long) {}
|
||||
|
||||
companion object {
|
||||
/** Shared no-op listener used as the default. */
|
||||
val None: RelayConnectionListener = object : RelayConnectionListener {}
|
||||
}
|
||||
}
|
||||
+17
@@ -47,11 +47,14 @@ import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.channels.ClosedSendChannelException
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlin.concurrent.atomics.AtomicLong
|
||||
import kotlin.concurrent.atomics.ExperimentalAtomicApi
|
||||
|
||||
/**
|
||||
* Represents an active session between a Nostr client and the relay.
|
||||
* Each one of these is a connection that can hold many subscriptions
|
||||
*/
|
||||
@OptIn(ExperimentalAtomicApi::class)
|
||||
class RelaySession(
|
||||
private val store: SessionBackend,
|
||||
val policy: IRelayPolicy,
|
||||
@@ -59,6 +62,13 @@ class RelaySession(
|
||||
private val onSend: (String) -> Unit,
|
||||
private val onClose: (RelaySession) -> Unit,
|
||||
negentropySettings: NegentropySettings = NegentropySettings.Default,
|
||||
/**
|
||||
* Stable, process-unique identifier for this connection. Used by the
|
||||
* server classes to key their connection registry and by
|
||||
* [RelayConnectionListener] callbacks so observers can correlate the
|
||||
* open/close of the same connection. Defaults to a fresh monotonic id.
|
||||
*/
|
||||
val id: Long = nextConnectionId(),
|
||||
) : AutoCloseable {
|
||||
private val subscriptions = LargeCache<String, Job>()
|
||||
|
||||
@@ -267,4 +277,11 @@ class RelaySession(
|
||||
init {
|
||||
policy.onConnect(::send)
|
||||
}
|
||||
|
||||
companion object {
|
||||
private val connectionIdSeq = AtomicLong(0L)
|
||||
|
||||
/** Allocates the next process-unique connection id. */
|
||||
fun nextConnectionId(): Long = connectionIdSeq.fetchAndAdd(1L)
|
||||
}
|
||||
}
|
||||
|
||||
+40
-16
@@ -26,6 +26,8 @@ import com.vitorpamplona.quartz.utils.cache.LargeCache
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlin.concurrent.atomics.AtomicLong
|
||||
import kotlin.concurrent.atomics.ExperimentalAtomicApi
|
||||
import kotlin.coroutines.CoroutineContext
|
||||
|
||||
/**
|
||||
@@ -62,20 +64,30 @@ import kotlin.coroutines.CoroutineContext
|
||||
* @param parentContext Parent coroutine context for all subscriptions.
|
||||
* @param negentropySettings NIP-77 tuning. Negentropy is effectively a no-op
|
||||
* here (the snapshot is empty) but the setting is plumbed for symmetry.
|
||||
* @param listener Observability hook fired as connections open and close,
|
||||
* keyed by [RelaySession.id]. Defaults to a no-op.
|
||||
*/
|
||||
@OptIn(ExperimentalAtomicApi::class)
|
||||
class ReqResponderServer(
|
||||
responder: ReqResponder,
|
||||
private val policyBuilder: () -> IRelayPolicy = { EmptyPolicy },
|
||||
parentContext: CoroutineContext = SupervisorJob(),
|
||||
private val negentropySettings: NegentropySettings = NegentropySettings.Default,
|
||||
private val listener: RelayConnectionListener = RelayConnectionListener.None,
|
||||
) : AutoCloseable {
|
||||
/** Scope for all subscriptions. */
|
||||
private val scope = CoroutineScope(parentContext + SupervisorJob())
|
||||
|
||||
private val backend = ReqResponderBackend(responder)
|
||||
|
||||
/** Active client sessions keyed by an opaque connection id. */
|
||||
private val connections = LargeCache<Int, RelaySession>()
|
||||
/** Live count of registered connections; backs [activeConnections]. */
|
||||
private val activeCount = AtomicLong(0L)
|
||||
|
||||
/** Active client sessions keyed by [RelaySession.id]. */
|
||||
private val connections = LargeCache<Long, RelaySession>()
|
||||
|
||||
/** Number of connections currently registered with this server. */
|
||||
val activeConnections: Long get() = activeCount.load()
|
||||
|
||||
/**
|
||||
* Registers a new client connection.
|
||||
@@ -83,19 +95,27 @@ class ReqResponderServer(
|
||||
* @param send Callback the server uses to send JSON messages to this client.
|
||||
* Implementations must be safe to call from any coroutine.
|
||||
*/
|
||||
fun connect(send: (String) -> Unit) =
|
||||
RelaySession(
|
||||
policy = policyBuilder(),
|
||||
store = backend,
|
||||
scope = scope,
|
||||
onSend = send,
|
||||
onClose = { session ->
|
||||
connections.remove(session.hashCode())
|
||||
},
|
||||
negentropySettings = negentropySettings,
|
||||
).also { session ->
|
||||
connections.put(session.hashCode(), session)
|
||||
}
|
||||
fun connect(send: (String) -> Unit): RelaySession {
|
||||
val session =
|
||||
RelaySession(
|
||||
policy = policyBuilder(),
|
||||
store = backend,
|
||||
scope = scope,
|
||||
onSend = send,
|
||||
onClose = { closed ->
|
||||
// Idempotent teardown accounting (see NostrServer.connect).
|
||||
if (connections.remove(closed.id) != null) {
|
||||
activeCount.addAndFetch(-1L)
|
||||
listener.onDisconnect(closed.id)
|
||||
}
|
||||
},
|
||||
negentropySettings = negentropySettings,
|
||||
)
|
||||
connections.put(session.id, session)
|
||||
activeCount.addAndFetch(1L)
|
||||
listener.onConnect(session.id)
|
||||
return session
|
||||
}
|
||||
|
||||
/**
|
||||
* Registers a new client connection and serves it for the duration of
|
||||
@@ -119,8 +139,12 @@ class ReqResponderServer(
|
||||
|
||||
/** Shuts down the server, cancelling all subscriptions. */
|
||||
override fun close() {
|
||||
connections.forEach { _, session -> session.cancelAllSubscriptions() }
|
||||
connections.forEach { _, session ->
|
||||
session.cancelAllSubscriptions()
|
||||
listener.onDisconnect(session.id)
|
||||
}
|
||||
connections.clear()
|
||||
activeCount.store(0L)
|
||||
scope.cancel()
|
||||
}
|
||||
}
|
||||
|
||||
+134
@@ -0,0 +1,134 @@
|
||||
/*
|
||||
* 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.server
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.policies.EmptyPolicy
|
||||
import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
|
||||
import kotlinx.coroutines.ExperimentalCoroutinesApi
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.emptyFlow
|
||||
import kotlinx.coroutines.test.UnconfinedTestDispatcher
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
@OptIn(ExperimentalCoroutinesApi::class)
|
||||
class RelayConnectionListenerTest {
|
||||
private class RecordingListener : RelayConnectionListener {
|
||||
val connected = mutableListOf<Long>()
|
||||
val disconnected = mutableListOf<Long>()
|
||||
|
||||
override fun onConnect(connectionId: Long) {
|
||||
connected.add(connectionId)
|
||||
}
|
||||
|
||||
override fun onDisconnect(connectionId: Long) {
|
||||
disconnected.add(connectionId)
|
||||
}
|
||||
}
|
||||
|
||||
private val emptyResponder =
|
||||
object : ReqResponder {
|
||||
override fun respond(filters: List<Filter>): Flow<Event> = emptyFlow()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun firesConnectAndDisconnectWithStableIds() =
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val listener = RecordingListener()
|
||||
val server = ReqResponderServer(emptyResponder, parentContext = dispatcher, listener = listener)
|
||||
|
||||
val a = server.connect {}
|
||||
val b = server.connect {}
|
||||
|
||||
assertEquals(listOf(a.id, b.id), listener.connected)
|
||||
assertTrue(a.id != b.id, "connection ids must be unique")
|
||||
assertEquals(2L, server.activeConnections)
|
||||
|
||||
a.close()
|
||||
assertEquals(listOf(a.id), listener.disconnected)
|
||||
assertEquals(1L, server.activeConnections)
|
||||
|
||||
b.close()
|
||||
assertEquals(listOf(a.id, b.id), listener.disconnected)
|
||||
assertEquals(0L, server.activeConnections)
|
||||
|
||||
server.close()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun doubleCloseIsAccountedOnce() =
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val listener = RecordingListener()
|
||||
val server = ReqResponderServer(emptyResponder, parentContext = dispatcher, listener = listener)
|
||||
|
||||
val s = server.connect {}
|
||||
s.close()
|
||||
s.close()
|
||||
|
||||
assertEquals(1, listener.disconnected.size)
|
||||
assertEquals(0L, server.activeConnections)
|
||||
|
||||
server.close()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun serverCloseDisconnectsRemaining() =
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val listener = RecordingListener()
|
||||
val server = ReqResponderServer(emptyResponder, parentContext = dispatcher, listener = listener)
|
||||
|
||||
val s = server.connect {}
|
||||
assertEquals(1L, server.activeConnections)
|
||||
|
||||
server.close()
|
||||
|
||||
assertEquals(listOf(s.id), listener.disconnected)
|
||||
assertEquals(0L, server.activeConnections)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun nostrServerReportsActiveConnections() =
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val listener = RecordingListener()
|
||||
NostrServer(
|
||||
store = EventStore(null),
|
||||
policyBuilder = { EmptyPolicy },
|
||||
parentContext = dispatcher,
|
||||
listener = listener,
|
||||
).use { server ->
|
||||
val s = server.connect {}
|
||||
assertEquals(1L, server.activeConnections)
|
||||
assertEquals(listOf(s.id), listener.connected)
|
||||
|
||||
s.close()
|
||||
assertEquals(0L, server.activeConnections)
|
||||
assertEquals(listOf(s.id), listener.disconnected)
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user