mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
refactor(quartz): simplify relay server internals (no behavior change)
Cleanup pass (4 review angles), behaviour-preserving, full suite green:
- Extract RelayServerBase: NostrServer and ReqResponderServer duplicated
connect/serve/buildPolicy/activeConnections and the connection scope verbatim
(ConnectionRegistry only unified the bookkeeping). They now share one engine
base and contribute only their backend + teardown; ReqResponderServer is ~10
lines, NostrServer just its ingest/store wiring.
- HyperLogLog.leadingZeroBits: replace the hand-rolled per-byte bit loop with
stdlib Int.countLeadingZeroBits() (- 24 for the 0..255 byte).
- LimitsPolicy.clampLimits: drop the `var changed` + throwaway-list map for a
`none{} -> map{}` that's simpler and allocates nothing when no filter is
clamped (the common case once maxLimit is set).
- PolicyStack: collapse the two first-non-null hook loops to firstNotNullOfOrNull.
- RelaySession.handleAuth: use OkMessage.rejected(MachineReadablePrefix.ERROR, …)
instead of hand-writing the "error:" prefix, matching the COUNT path.
Skipped: rewriting FullAuthPolicy's pre-existing auth-required:/invalid: reason
strings to MachineReadablePrefix — unchanged context lines, out of this diff's
scope.
https://claude.ai/code/session_016YBS2pWCBSDgAMthHzfCTr
This commit is contained in:
+13
-72
@@ -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()
|
||||
|
||||
+123
@@ -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()
|
||||
}
|
||||
}
|
||||
+1
-1
@@ -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
|
||||
}
|
||||
|
||||
|
||||
+5
-65
@@ -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)
|
||||
}
|
||||
|
||||
+3
-13
@@ -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<Filter>): List<Filter> {
|
||||
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? =
|
||||
|
||||
+2
-12
@@ -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 <T : Command> runPolicies(
|
||||
initialCmd: T,
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user