diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/NostrServer.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/NostrServer.kt index adb3ee61b3..f7e75ef1a1 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/NostrServer.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/NostrServer.kt @@ -21,20 +21,17 @@ package com.vitorpamplona.quartz.nip01Core.relay.server import com.vitorpamplona.quartz.nip01Core.crypto.verify -import com.vitorpamplona.quartz.nip01Core.relay.server.policies.LimitsPolicy import com.vitorpamplona.quartz.nip01Core.relay.server.policies.VerifyPolicy import com.vitorpamplona.quartz.nip01Core.store.IEventStore import com.vitorpamplona.quartz.nip77Negentropy.NegentropySettings -import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.cancel import kotlin.coroutines.CoroutineContext /** - * Represents a Nostr relay server that manages client connections, event storage, and verification. - * - * This class acts as the central coordinator for a relay server, handling the lifecycle of [RelaySession]s - * and providing access to the underlying event store. + * Storage-backed Nostr relay server: a [RelayServerBase] whose [backend] is a + * [LiveEventStore] over [store] (historical replay + live tail after EOSE), + * fed through a group-commit [IngestQueue]. * * @param store The [IEventStore] backing this relay. * @param policyBuilder Controls requirements for relay commands. @@ -50,22 +47,19 @@ import kotlin.coroutines.CoroutineContext * @param listener Observability hook fired as connections open and close, * keyed by [RelaySession.id]. Defaults to a no-op. * @param limits Operational limits enforced on every connection (per-command - * via a composed [LimitsPolicy], plus the session-level message-size and - * subscription caps) and advertised via [RelayLimits.toNip11Limitation]. - * Null disables limit enforcement. + * via a composed [com.vitorpamplona.quartz.nip01Core.relay.server.policies.LimitsPolicy], + * plus the session-level message-size and subscription caps) and advertised + * via [RelayLimits.toNip11Limitation]. Null disables limit enforcement. */ class NostrServer( private val store: IEventStore, - private val policyBuilder: () -> IRelayPolicy = { VerifyPolicy }, - private val parentContext: CoroutineContext = SupervisorJob(), + policyBuilder: () -> IRelayPolicy = { VerifyPolicy }, + parentContext: CoroutineContext = SupervisorJob(), parallelVerify: Boolean = false, - private val negentropySettings: NegentropySettings = NegentropySettings.Default, + negentropySettings: NegentropySettings = NegentropySettings.Default, listener: RelayServerListener = RelayServerListener.None, - val limits: RelayLimits? = null, -) : AutoCloseable { - /** Scope for all subscriptions. */ - private val scope = CoroutineScope(parentContext + SupervisorJob()) - + limits: RelayLimits? = null, +) : RelayServerBase(policyBuilder, parentContext, negentropySettings, listener, limits) { /** * Group-commit writer shared across every connected session. * Sessions hand off EVENT publishes here instead of awaiting @@ -80,66 +74,13 @@ class NostrServer( verify = if (parallelVerify) ({ it.verify() }) else null, ) - private val subStore = LiveEventStore(store, ingest) - - private val connections = ConnectionRegistry(listener) - - /** Number of connections currently registered with this server. */ - val activeConnections: Long get() = connections.active - - /** - * Builds the per-connection policy, prepending a [LimitsPolicy] when - * [limits] declares per-command caps so requests are clamped/rejected - * before the application policy runs. - */ - private fun buildPolicy(): IRelayPolicy { - val base = policyBuilder() - return if (limits != null) LimitsPolicy(limits) + base else base - } - - /** - * Registers a new client connection. - * - * @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 = - connections.register( - RelaySession( - policy = buildPolicy(), - store = subStore, - scope = scope, - onSend = send, - onClose = { connections.unregister(it.id) }, - negentropySettings = negentropySettings, - ), - ) - - /** - * Registers a new client connection and serves it for the duration of - * [incoming]. The session is automatically closed when [incoming] returns. - * - * @param send Callback the server uses to send JSON messages to this client. - * @param incoming Suspend block that yields raw JSON strings from the client - * (e.g., reading WebSocket text frames in a loop). - */ - suspend fun serve( - send: (String) -> Unit, - incoming: suspend (RelaySession) -> Unit, - ) { - val session = connect(send) - try { - incoming(session) - } finally { - session.close() - } - } + override val backend: SessionBackend = LiveEventStore(store, ingest) /** * Shuts down the server, cancelling all subscriptions and closing the store. */ override fun close() { - connections.closeAll() + closeConnections() ingest.close() scope.cancel() store.close() diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelayServerBase.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelayServerBase.kt new file mode 100644 index 0000000000..a544c74dbc --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelayServerBase.kt @@ -0,0 +1,123 @@ +/* + * 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.relay.server.policies.LimitsPolicy +import com.vitorpamplona.quartz.nip77Negentropy.NegentropySettings +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.cancel +import kotlin.coroutines.CoroutineContext + +/** + * Shared engine behind [NostrServer] and [ReqResponderServer]. Owns the + * per-connection coroutine [scope], the [ConnectionRegistry] (and the + * [activeConnections] gauge / [RelayServerListener] dispatch it drives), the + * policy composition, and the connect/serve plumbing — so the two concrete + * servers differ only in the [backend] they wrap and what their [close] tears + * down beyond the connections. + * + * @param policyBuilder Builds a fresh base policy per connection. + * @param parentContext Parent coroutine context for all subscriptions. + * @param negentropySettings NIP-77 tuning handed to each session. + * @param listener Connection-lifecycle observer. + * @param limits Operational limits; when non-null a [LimitsPolicy] is prepended + * to every connection's policy and the value is exposed for NIP-11 advertising. + */ +abstract class RelayServerBase( + private val policyBuilder: () -> IRelayPolicy, + parentContext: CoroutineContext, + private val negentropySettings: NegentropySettings, + listener: RelayServerListener, + val limits: RelayLimits?, +) : AutoCloseable { + /** Scope for all subscriptions. */ + protected val scope = CoroutineScope(parentContext + SupervisorJob()) + + private val connections = ConnectionRegistry(listener) + + /** The data plane every connection on this server talks to. */ + protected abstract val backend: SessionBackend + + /** Number of connections currently registered with this server. */ + val activeConnections: Long get() = connections.active + + /** + * Builds the per-connection policy, prepending a [LimitsPolicy] when + * [limits] is set so requests are clamped/rejected before the application + * policy runs. + */ + private fun buildPolicy(): IRelayPolicy { + val base = policyBuilder() + return if (limits != null) LimitsPolicy(limits) + base else base + } + + /** + * Registers a new client connection. + * + * @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 = + connections.register( + RelaySession( + policy = buildPolicy(), + store = backend, + scope = scope, + onSend = send, + onClose = { connections.unregister(it.id) }, + negentropySettings = negentropySettings, + ), + ) + + /** + * Registers a new client connection and serves it for the duration of + * [incoming]. The session is automatically closed when [incoming] returns. + * + * @param send Callback the server uses to send JSON messages to this client. + * @param incoming Suspend block that yields raw JSON strings from the client + * (e.g., reading WebSocket text frames in a loop). + */ + suspend fun serve( + send: (String) -> Unit, + incoming: suspend (RelaySession) -> Unit, + ) { + val session = connect(send) + try { + incoming(session) + } finally { + session.close() + } + } + + /** Cancels every open connection's subscriptions and clears the registry. */ + protected fun closeConnections() = connections.closeAll() + + /** + * Shuts the server down: cancels all subscriptions and the server [scope]. + * Subclasses override to tear down extra resources (e.g. an event store), + * calling [closeConnections] and cancelling [scope] in the order they need. + */ + override fun close() { + closeConnections() + scope.cancel() + } +} diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelaySession.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelaySession.kt index f480ef4b6e..7502c3969d 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelaySession.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelaySession.kt @@ -214,7 +214,7 @@ class RelaySession( } catch (e: CancellationException) { throw e } catch (e: Exception) { - send(OkMessage(cmd.event.id, false, "error: ${e.message ?: "authentication failed"}")) + send(OkMessage.rejected(cmd.event.id, MachineReadablePrefix.ERROR, e.message ?: "authentication failed")) return } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/ReqResponderServer.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/ReqResponderServer.kt index c28d62f5db..85a02417b0 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/ReqResponderServer.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/ReqResponderServer.kt @@ -21,11 +21,8 @@ package com.vitorpamplona.quartz.nip01Core.relay.server import com.vitorpamplona.quartz.nip01Core.relay.server.policies.EmptyPolicy -import com.vitorpamplona.quartz.nip01Core.relay.server.policies.LimitsPolicy import com.vitorpamplona.quartz.nip77Negentropy.NegentropySettings -import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.SupervisorJob -import kotlinx.coroutines.cancel import kotlin.coroutines.CoroutineContext /** @@ -71,68 +68,11 @@ import kotlin.coroutines.CoroutineContext */ class ReqResponderServer( responder: ReqResponder, - private val policyBuilder: () -> IRelayPolicy = { EmptyPolicy }, + policyBuilder: () -> IRelayPolicy = { EmptyPolicy }, parentContext: CoroutineContext = SupervisorJob(), - private val negentropySettings: NegentropySettings = NegentropySettings.Default, + negentropySettings: NegentropySettings = NegentropySettings.Default, listener: RelayServerListener = RelayServerListener.None, - val limits: RelayLimits? = null, -) : AutoCloseable { - /** Scope for all subscriptions. */ - private val scope = CoroutineScope(parentContext + SupervisorJob()) - - private val backend = ReqResponderBackend(responder) - - private val connections = ConnectionRegistry(listener) - - /** Number of connections currently registered with this server. */ - val activeConnections: Long get() = connections.active - - private fun buildPolicy(): IRelayPolicy { - val base = policyBuilder() - return if (limits != null) LimitsPolicy(limits) + base else base - } - - /** - * Registers a new client connection. - * - * @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 = - connections.register( - RelaySession( - policy = buildPolicy(), - store = backend, - scope = scope, - onSend = send, - onClose = { connections.unregister(it.id) }, - negentropySettings = negentropySettings, - ), - ) - - /** - * Registers a new client connection and serves it for the duration of - * [incoming]. The session is automatically closed when [incoming] returns. - * - * @param send Callback the server uses to send JSON messages to this client. - * @param incoming Suspend block that yields raw JSON strings from the client - * (e.g., reading WebSocket text frames in a loop). - */ - suspend fun serve( - send: (String) -> Unit, - incoming: suspend (RelaySession) -> Unit, - ) { - val session = connect(send) - try { - incoming(session) - } finally { - session.close() - } - } - - /** Shuts down the server, cancelling all subscriptions. */ - override fun close() { - connections.closeAll() - scope.cancel() - } + limits: RelayLimits? = null, +) : RelayServerBase(policyBuilder, parentContext, negentropySettings, listener, limits) { + override val backend: SessionBackend = ReqResponderBackend(responder) } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/LimitsPolicy.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/LimitsPolicy.kt index d413f59566..724a6e6c4a 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/LimitsPolicy.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/LimitsPolicy.kt @@ -101,21 +101,11 @@ class LimitsPolicy( private fun invalid(message: String): String = MachineReadablePrefix.INVALID.format(message) - /** Returns the same list reference when nothing changes, so callers can skip rebuilding the command. */ + /** Returns the same list reference (no allocation) when no filter needs clamping. */ private fun clampLimits(filters: List): List { if (limits.maxLimit == null && limits.defaultLimit == null) return filters - var changed = false - val out = - filters.map { f -> - val target = targetLimit(f.limit) - if (target != f.limit) { - changed = true - f.copy(limit = target) - } else { - f - } - } - return if (changed) out else filters + if (filters.none { targetLimit(it.limit) != it.limit }) return filters + return filters.map { it.copy(limit = targetLimit(it.limit)) } } private fun targetLimit(current: Int?): Int? = diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/PolicyStack.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/PolicyStack.kt index c0654a4ffd..1a5687c996 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/PolicyStack.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/PolicyStack.kt @@ -56,22 +56,12 @@ class PolicyStack( policies.forEach { it.onAuthenticated(pubKey, event) } } - override fun acceptMessage(message: String): String? { - for (policy in policies) { - policy.acceptMessage(message)?.let { return it } - } - return null - } + override fun acceptMessage(message: String): String? = policies.firstNotNullOfOrNull { it.acceptMessage(message) } override fun acceptSubscription( subId: String, openSubscriptions: Int, - ): String? { - for (policy in policies) { - policy.acceptSubscription(subId, openSubscriptions)?.let { return it } - } - return null - } + ): String? = policies.firstNotNullOfOrNull { it.acceptSubscription(subId, openSubscriptions) } private inline fun runPolicies( initialCmd: T, diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip45Count/HyperLogLog.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip45Count/HyperLogLog.kt index e00ca410bc..3cfa916dc9 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip45Count/HyperLogLog.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip45Count/HyperLogLog.kt @@ -138,7 +138,8 @@ object HyperLogLog { /** * Counts consecutive zero bits (MSB first) from byte [fromByteIndex] to the - * end of [bytes]. KMP-safe (no `Integer.numberOfLeadingZeros`). + * end of [bytes]. A byte value in 0..255 occupies the low 8 bits of an Int, + * so its in-byte leading zeros are `countLeadingZeroBits() - 24`. */ private fun leadingZeroBits( bytes: ByteArray, @@ -150,12 +151,7 @@ object HyperLogLog { if (b == 0) { count += 8 } else { - var mask = 0x80 - while (mask != 0 && (b and mask) == 0) { - count++ - mask = mask shr 1 - } - return count + return count + b.countLeadingZeroBits() - 24 } } return count