mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
refactor(quartz): rename ReqResponder -> EventSource; drop JwtAuthPolicy doc
The SPI name leaned on "Req" (the NIP-01 command), which isn't self-explanatory. Rename to read as what it is — a source of events for a query: - ReqResponder -> EventSource (method respond() -> events()) - ReqResponderBackend -> EventSourceBackend (param responder -> source) - ReqResponderServer -> EventSourceServer - *Test + RELAY.md + KDoc references updated to match. Also drop the JwtAuthPolicy snippet from RELAY.md (it was only an illustrative doc example, never a real class); the surrounding prose still documents the `authorize` bridge hook. Pure rename + doc edit, no behavior change; full suite green. https://claude.ai/code/session_016YBS2pWCBSDgAMthHzfCTr
This commit is contained in:
+19
-30
@@ -2,7 +2,7 @@
|
||||
|
||||
Quartz provides a transport-agnostic relay engine. You provide a `send` callback per connection, and it gives you a `RelaySession` that accepts raw JSON strings. Plug it into Ktor or any WebSocket transport.
|
||||
|
||||
There are two engines on the same `RelaySession` core: `NostrServer` for storage-backed relays (an `IEventStore` with a live tail after EOSE) and `ReqResponderServer` for non-storage relays (search, redirector, computed data — see [Non-Storage Relays](#non-storage-relays-search-redirector-computed)).
|
||||
There are two engines on the same `RelaySession` core: `NostrServer` for storage-backed relays (an `IEventStore` with a live tail after EOSE) and `EventSourceServer` for non-storage relays (search, redirector, computed data — see [Non-Storage Relays](#non-storage-relays-search-redirector-computed)).
|
||||
|
||||
Both `NostrServer` and `EventStore` implement `AutoCloseable`.
|
||||
|
||||
@@ -86,14 +86,14 @@ By default, all single-letter tags with values are indexed. Override `shouldInde
|
||||
|
||||
A relay's job is to *answer REQs*. When the answer doesn't come from a stored
|
||||
set — a NIP-50 search that forwards to an HTTP backend, a relay that emits
|
||||
computed/projected data — implement `ReqResponder` and serve it with
|
||||
`ReqResponderServer`. You supply a `Flow<Event>`; the engine owns the wire
|
||||
computed/projected data — implement `EventSource` and serve it with
|
||||
`EventSourceServer`. You supply a `Flow<Event>`; the engine owns the wire
|
||||
protocol (challenge/auth, command parsing, policy, `EVENT`/`EOSE`/`CLOSED`
|
||||
framing, subscription lifecycle), so there's no hand-written read loop.
|
||||
|
||||
```kotlin
|
||||
class SearchResponder(private val backend: SearchApi) : ReqResponder {
|
||||
override fun respond(filters: List<Filter>): Flow<Event> = flow {
|
||||
class SearchEventSource(private val backend: SearchApi) : EventSource {
|
||||
override fun events(filters: List<Filter>): Flow<Event> = flow {
|
||||
filters.forEach { f ->
|
||||
f.search?.let { raw ->
|
||||
val q = SearchQuery.parse(raw)
|
||||
@@ -102,12 +102,12 @@ class SearchResponder(private val backend: SearchApi) : ReqResponder {
|
||||
}
|
||||
}
|
||||
}
|
||||
// COUNT defaults to counting respond(); override for a cheaper backend count.
|
||||
// COUNT defaults to counting events(); override for a cheaper backend count.
|
||||
}
|
||||
|
||||
fun main() {
|
||||
val server = ReqResponderServer(
|
||||
responder = SearchResponder(searchApi),
|
||||
val server = EventSourceServer(
|
||||
source = SearchEventSource(searchApi),
|
||||
policyBuilder = { FullAuthPolicy(relay) }, // optional NIP-42 gating
|
||||
)
|
||||
|
||||
@@ -126,7 +126,7 @@ fun main() {
|
||||
}
|
||||
```
|
||||
|
||||
**EOSE = flow completion.** The engine sends `EOSE` when `respond(...)`
|
||||
**EOSE = flow completion.** The engine sends `EOSE` when `events(...)`
|
||||
completes, which is the natural shape for finite queries. Relays that need an
|
||||
open-ended live tail after EOSE should use the storage path (`NostrServer` +
|
||||
`IEventStore`) instead. EVENT publishes are rejected (`OK false` —
|
||||
@@ -134,10 +134,10 @@ open-ended live tail after EOSE should use the storage path (`NostrServer` +
|
||||
there's no stored set. A failure thrown from the flow ends the subscription
|
||||
with `CLOSED` `error: <message>`.
|
||||
|
||||
`ReqResponderServer` and `NostrServer` are the two concrete dispatch engines;
|
||||
`EventSourceServer` and `NostrServer` are the two concrete dispatch engines;
|
||||
both build on `RelaySession` and the shared `SessionBackend` seam (`LiveEventStore`
|
||||
is the storage-backed `SessionBackend`; `ReqResponderBackend` adapts a
|
||||
`ReqResponder`). Implement `SessionBackend` directly only if you need custom
|
||||
is the storage-backed `SessionBackend`; `EventSourceBackend` adapts a
|
||||
`EventSource`). Implement `SessionBackend` directly only if you need custom
|
||||
control over the EVENT/negentropy paths as well as REQ/COUNT.
|
||||
|
||||
## Policies
|
||||
@@ -169,19 +169,7 @@ val server = NostrServer(
|
||||
)
|
||||
```
|
||||
|
||||
`FullAuthPolicy` already implements the full NIP-42 challenge/verify handshake — you should not re-implement it. To bridge auth to an external system (e.g. exchange the verified event for a backend JWT), override the `suspend` `authorize` hook. It runs after the NIP-42 checks pass (and after the rest of the policy chain approves), and the pubkey is recorded only once it returns — so it can do network/disk I/O, and throwing from it rejects the login (`OK false`) with the connection left unauthenticated (nothing was committed):
|
||||
|
||||
```kotlin
|
||||
class JwtAuthPolicy(
|
||||
relay: NormalizedRelayUrl,
|
||||
private val backend: AuthBackend,
|
||||
) : FullAuthPolicy(relay) {
|
||||
override suspend fun authorize(pubKey: HexKey, event: RelayAuthEvent) {
|
||||
// Suspends; a thrown exception rejects the login with OK false.
|
||||
backend.exchangeForSession(pubKey, event)
|
||||
}
|
||||
}
|
||||
```
|
||||
`FullAuthPolicy` already implements the full NIP-42 challenge/verify handshake — you should not re-implement it. To bridge auth to an external system (e.g. exchange the verified event for a backend JWT), override the `suspend` `authorize` hook. It runs after the NIP-42 checks pass (and after the rest of the policy chain approves), and the pubkey is recorded only once it returns — so it can do network/disk I/O, and throwing from it rejects the login (`OK false`) with the connection left unauthenticated (nothing was committed).
|
||||
|
||||
### Composing Policies
|
||||
|
||||
@@ -373,7 +361,7 @@ are omitted), and `CONTENT_TYPE` is `application/nostr+json`.
|
||||
## Approximate COUNT (NIP-45 HyperLogLog)
|
||||
|
||||
`COUNT` is answered by `SessionBackend.countResult(filters)` (and
|
||||
`ReqResponder.countResult`), which defaults to an exact count. To return a
|
||||
`EventSource.countResult`), which defaults to an exact count. To return a
|
||||
mergeable HyperLogLog estimate instead — for the six canonical NIP-45 queries
|
||||
(reaction/repost/quote/reply/comment/follower counts) — fold matching pubkeys
|
||||
into an `HllBuilder` and return its `CountResult`:
|
||||
@@ -421,10 +409,11 @@ coroutine — keep them cheap and non-blocking.
|
||||
```
|
||||
quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/
|
||||
├── relay/server/
|
||||
│ ├── RelayServerBase.kt # Shared engine (connect/serve/policy/registry)
|
||||
│ ├── NostrServer.kt # Storage-backed engine (IEventStore)
|
||||
│ ├── ReqResponderServer.kt # Non-storage engine (search/redirector/computed)
|
||||
│ ├── ReqResponder.kt # Flow<Event> REQ-responder SPI
|
||||
│ ├── SessionBackend.kt # Data-plane seam (LiveEventStore / ReqResponderBackend)
|
||||
│ ├── EventSourceServer.kt # Non-storage engine (search/redirector/computed)
|
||||
│ ├── EventSource.kt # Flow<Event> SPI for non-storage relays
|
||||
│ ├── SessionBackend.kt # Data-plane seam (LiveEventStore / EventSourceBackend)
|
||||
│ ├── RelayServerListener.kt # Connection open/close observability hook
|
||||
│ ├── RelaySession.kt # Per-connection handler (stable .id)
|
||||
│ ├── LiveEventStore.kt # Reactive event streaming (storage SessionBackend)
|
||||
@@ -433,7 +422,7 @@ quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/
|
||||
│ └── policies/
|
||||
│ ├── EmptyPolicy.kt # Accept everything
|
||||
│ ├── VerifyPolicy.kt # Signature verification (default)
|
||||
│ ├── FullAuthPolicy.kt # NIP-42 auth required (override onAuthenticated to bridge)
|
||||
│ ├── FullAuthPolicy.kt # NIP-42 auth required (override authorize to bridge)
|
||||
│ ├── LimitsPolicy.kt # Per-command enforcement of RelayLimits
|
||||
│ └── PolicyStack.kt # Chain multiple policies
|
||||
├── relay/commands/
|
||||
|
||||
+1
-1
@@ -27,7 +27,7 @@ import kotlin.concurrent.atomics.ExperimentalAtomicApi
|
||||
/**
|
||||
* Tracks the live [RelaySession]s of one relay server and fires the
|
||||
* [RelayServerListener] as they come and go. Shared by [NostrServer] and
|
||||
* [ReqResponderServer] so the connection bookkeeping — the stable-id keying,
|
||||
* [EventSourceServer] so the connection bookkeeping — the stable-id keying,
|
||||
* the [active] gauge, and the once-only teardown accounting — lives in exactly
|
||||
* one place.
|
||||
*/
|
||||
|
||||
+9
-9
@@ -31,13 +31,13 @@ import kotlinx.coroutines.flow.count
|
||||
* not backed by an event store — search relays, redirectors that forward to an
|
||||
* HTTP backend, and relays that emit computed/projected data.
|
||||
*
|
||||
* The dispatch engine ([ReqResponderServer]) owns the wire protocol — reading
|
||||
* The dispatch engine ([EventSourceServer]) owns the wire protocol — reading
|
||||
* frames, parsing commands, running the [IRelayPolicy], framing `EVENT`/`EOSE`/
|
||||
* `CLOSED`, and managing subscriptions. You only supply the events:
|
||||
*
|
||||
* ```
|
||||
* class SearchResponder(private val backend: SearchApi) : ReqResponder {
|
||||
* override fun respond(filters: List<Filter>): Flow<Event> = flow {
|
||||
* class SearchEventSource(private val backend: SearchApi) : EventSource {
|
||||
* override fun events(filters: List<Filter>): Flow<Event> = flow {
|
||||
* for (f in filters) {
|
||||
* f.search?.let { raw ->
|
||||
* val query = SearchQuery.parse(raw)
|
||||
@@ -50,7 +50,7 @@ import kotlinx.coroutines.flow.count
|
||||
*
|
||||
* ## EOSE semantics
|
||||
*
|
||||
* The engine sends `EOSE` when the returned [Flow] **completes**. A responder
|
||||
* The engine sends `EOSE` when the returned [Flow] **completes**. A source
|
||||
* therefore emits its full result set and then completes — the natural shape
|
||||
* for finite queries like search and counts. Relays that need an open-ended
|
||||
* live tail after EOSE should use the storage path ([NostrServer] with an
|
||||
@@ -58,22 +58,22 @@ import kotlinx.coroutines.flow.count
|
||||
*
|
||||
* Cancellation of the subscription (NIP-01 `CLOSE` or a dropped connection)
|
||||
* cancels collection of the flow, so respect coroutine cancellation in
|
||||
* [respond].
|
||||
* [events].
|
||||
*/
|
||||
interface ReqResponder {
|
||||
interface EventSource {
|
||||
/**
|
||||
* Returns the events matching [filters]. Emit each match and then let the
|
||||
* flow complete; completion is what triggers `EOSE`. The [filters] are the
|
||||
* (possibly policy-rewritten) filters from the REQ.
|
||||
*/
|
||||
fun respond(filters: List<Filter>): Flow<Event>
|
||||
fun events(filters: List<Filter>): Flow<Event>
|
||||
|
||||
/**
|
||||
* Answers a NIP-45 COUNT. The default counts the events produced by
|
||||
* [respond]; override it when the backend can count without materializing
|
||||
* [events]; override it when the backend can count without materializing
|
||||
* every event.
|
||||
*/
|
||||
suspend fun count(filters: List<Filter>): Int = respond(filters).count()
|
||||
suspend fun count(filters: List<Filter>): Int = events(filters).count()
|
||||
|
||||
/**
|
||||
* Answers a NIP-45 COUNT, optionally approximate and/or carrying a
|
||||
+8
-8
@@ -26,25 +26,25 @@ import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import kotlinx.coroutines.flow.collect
|
||||
|
||||
/**
|
||||
* Adapts a [ReqResponder] (the `Flow<Event>` SPI) into the [SessionBackend] the
|
||||
* dispatch engine consumes. Drains the responder's flow into the per-event
|
||||
* Adapts an [EventSource] (the `Flow<Event>` SPI) into the [SessionBackend] the
|
||||
* dispatch engine consumes. Drains the source's flow into the per-event
|
||||
* callback and signals EOSE on completion. The EVENT and negentropy paths use
|
||||
* the [SessionBackend] defaults (reject / empty), so a responder relay neither
|
||||
* the [SessionBackend] defaults (reject / empty), so a source relay neither
|
||||
* stores events nor reconciles.
|
||||
*/
|
||||
class ReqResponderBackend(
|
||||
private val responder: ReqResponder,
|
||||
class EventSourceBackend(
|
||||
private val source: EventSource,
|
||||
) : SessionBackend {
|
||||
override suspend fun query(
|
||||
filters: List<Filter>,
|
||||
onEach: (Event) -> Unit,
|
||||
onEose: () -> Unit,
|
||||
) {
|
||||
responder.respond(filters).collect { onEach(it) }
|
||||
source.events(filters).collect { onEach(it) }
|
||||
onEose()
|
||||
}
|
||||
|
||||
override suspend fun count(filters: List<Filter>): Int = responder.count(filters)
|
||||
override suspend fun count(filters: List<Filter>): Int = source.count(filters)
|
||||
|
||||
override suspend fun countResult(filters: List<Filter>): CountResult = responder.countResult(filters)
|
||||
override suspend fun countResult(filters: List<Filter>): CountResult = source.countResult(filters)
|
||||
}
|
||||
+8
-8
@@ -27,7 +27,7 @@ import kotlin.coroutines.CoroutineContext
|
||||
|
||||
/**
|
||||
* A transport-agnostic relay engine for relays that answer REQs from a
|
||||
* [ReqResponder] instead of an event store — search relays, redirectors that
|
||||
* [EventSource] instead of an event store — search relays, redirectors that
|
||||
* forward to an HTTP backend, and relays that emit computed/projected data.
|
||||
*
|
||||
* This is the storage-free sibling of [NostrServer]. It owns the full NIP-01
|
||||
@@ -35,11 +35,11 @@ import kotlin.coroutines.CoroutineContext
|
||||
* EOSE/CLOSED framing, and subscription lifecycle — so callers only provide a
|
||||
* per-connection `send` callback and feed it raw JSON frames, exactly like
|
||||
* [NostrServer]. EVENT publishes are rejected (there is nothing to store) and
|
||||
* negentropy is disabled, per [ReqResponderBackend] / [SessionBackend].
|
||||
* negentropy is disabled, per [EventSourceBackend] / [SessionBackend].
|
||||
*
|
||||
* ```
|
||||
* val server = ReqResponderServer(
|
||||
* responder = SearchResponder(searchApi),
|
||||
* val server = EventSourceServer(
|
||||
* source = SearchEventSource(searchApi),
|
||||
* policyBuilder = { FullAuthPolicy(relay) }, // optional NIP-42 gating
|
||||
* )
|
||||
*
|
||||
@@ -51,7 +51,7 @@ import kotlin.coroutines.CoroutineContext
|
||||
*
|
||||
* Both this class and [RelaySession] implement [AutoCloseable].
|
||||
*
|
||||
* @param responder Produces the events that answer each REQ.
|
||||
* @param source Produces the events that answer each REQ.
|
||||
* @param policyBuilder Builds a fresh [IRelayPolicy] per connection. Defaults to
|
||||
* [EmptyPolicy] (accept REQ/COUNT, no signature verification — appropriate for
|
||||
* a read-only relay). Pass a [com.vitorpamplona.quartz.nip01Core.relay.server.policies.FullAuthPolicy]
|
||||
@@ -66,13 +66,13 @@ import kotlin.coroutines.CoroutineContext
|
||||
* subscription caps) and advertised via [RelayLimits.toNip11Limitation].
|
||||
* Null disables limit enforcement.
|
||||
*/
|
||||
class ReqResponderServer(
|
||||
responder: ReqResponder,
|
||||
class EventSourceServer(
|
||||
source: EventSource,
|
||||
policyBuilder: () -> IRelayPolicy = { EmptyPolicy },
|
||||
parentContext: CoroutineContext = SupervisorJob(),
|
||||
negentropySettings: NegentropySettings = NegentropySettings.Default,
|
||||
listener: RelayServerListener = RelayServerListener.None,
|
||||
limits: RelayLimits? = null,
|
||||
) : RelayServerBase(policyBuilder, parentContext, negentropySettings, listener, limits) {
|
||||
override val backend: SessionBackend = ReqResponderBackend(responder)
|
||||
override val backend: SessionBackend = EventSourceBackend(source)
|
||||
}
|
||||
+1
-1
@@ -26,7 +26,7 @@ import com.vitorpamplona.quartz.nip11RelayInfo.Nip11RelayInformation.RelayInform
|
||||
* The relay's operational limits — the single source of truth that is both
|
||||
* **enforced** and **advertised**.
|
||||
*
|
||||
* Hand this to [NostrServer] / [ReqResponderServer] and it is applied to every
|
||||
* Hand this to [NostrServer] / [EventSourceServer] and it is applied to every
|
||||
* connection (per-command checks via [LimitsPolicy], plus the session-level
|
||||
* message-size and subscription-count caps in [RelaySession]); call
|
||||
* [toNip11Limitation] to advertise the very same numbers in the relay's NIP-11
|
||||
|
||||
+1
-1
@@ -28,7 +28,7 @@ import kotlinx.coroutines.cancel
|
||||
import kotlin.coroutines.CoroutineContext
|
||||
|
||||
/**
|
||||
* Shared engine behind [NostrServer] and [ReqResponderServer]. Owns the
|
||||
* Shared engine behind [NostrServer] and [EventSourceServer]. 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
|
||||
|
||||
+1
-1
@@ -22,7 +22,7 @@ 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
|
||||
* [NostrServer] and [EventSourceServer] 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
|
||||
|
||||
+1
-1
@@ -263,7 +263,7 @@ class RelaySession(
|
||||
// Subscription was closed – this is expected.
|
||||
throw e
|
||||
} catch (e: Exception) {
|
||||
// A backend failure (e.g. a responder's network I/O)
|
||||
// A backend failure (e.g. an event source's network I/O)
|
||||
// ends the subscription with a machine-readable CLOSED
|
||||
// rather than silently dropping the coroutine.
|
||||
send(ClosedMessage.of(cmd.subId, MachineReadablePrefix.ERROR, e.message ?: "query failed"))
|
||||
|
||||
+4
-4
@@ -35,10 +35,10 @@ import com.vitorpamplona.quartz.nip01Core.store.IdAndTime
|
||||
*
|
||||
* - [LiveEventStore] — the storage-backed path: replays stored events, signals
|
||||
* EOSE, then keeps streaming live inserts. Used by [NostrServer].
|
||||
* - [ReqResponderBackend] — adapts a [ReqResponder] for non-storage relays
|
||||
* (search, redirector, computed/projected data). Used by [ReqResponderServer].
|
||||
* - [EventSourceBackend] — adapts a [EventSource] for non-storage relays
|
||||
* (search, redirector, computed/projected data). Used by [EventSourceServer].
|
||||
*
|
||||
* Most non-storage relays should implement the higher-level [ReqResponder]
|
||||
* Most non-storage relays should implement the higher-level [EventSource]
|
||||
* (a `Flow<Event>` SPI) rather than this interface directly. Implement
|
||||
* [SessionBackend] only when you need control over the write path or negentropy
|
||||
* snapshots; [query] and [count] are the only members without a default.
|
||||
@@ -48,7 +48,7 @@ interface SessionBackend {
|
||||
* Answers a REQ. Calls [onEach] for every matching event, then [onEose]
|
||||
* once the stored set is exhausted. A storage backend keeps suspending
|
||||
* after [onEose] to stream live events until the subscription is
|
||||
* cancelled; a finite responder returns after [onEose].
|
||||
* cancelled; a finite source returns after [onEose].
|
||||
*/
|
||||
suspend fun query(
|
||||
filters: List<Filter>,
|
||||
|
||||
+26
-26
@@ -42,7 +42,7 @@ import kotlin.test.assertEquals
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
@OptIn(ExperimentalCoroutinesApi::class)
|
||||
class ReqResponderServerTest {
|
||||
class EventSourceServerTest {
|
||||
private val pubkey = "46fcbe3065eaf1ae7811465924e48923363ff3f526bd6f73d7c184b16bd8ce4d"
|
||||
private val sig = "4aa5264965018fa12a326686ad3d3bd8beae3218dcc83689b19ca1e6baeb791531943c15363aa6707c7c0c8b2d601deca1f20c32078b2872d356cdca03b04cce"
|
||||
|
||||
@@ -65,19 +65,19 @@ class ReqResponderServerTest {
|
||||
fun containing(label: String) = messages.filter { it.contains("\"$label\"") }
|
||||
}
|
||||
|
||||
/** A responder that returns a fixed list of events for any REQ. */
|
||||
private class FixedResponder(
|
||||
/** A source that returns a fixed list of events for any REQ. */
|
||||
private class FixedSource(
|
||||
private val events: List<Event>,
|
||||
) : ReqResponder {
|
||||
override fun respond(filters: List<Filter>): Flow<Event> = flowOf(*events.toTypedArray())
|
||||
) : EventSource {
|
||||
override fun events(filters: List<Filter>): Flow<Event> = flowOf(*events.toTypedArray())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun reqStreamsEventsThenEose() =
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val responder = FixedResponder(listOf(event(1), event(2)))
|
||||
ReqResponderServer(responder, parentContext = dispatcher).use { server ->
|
||||
val source = FixedSource(listOf(event(1), event(2)))
|
||||
EventSourceServer(source, parentContext = dispatcher).use { server ->
|
||||
val collector = MessageCollector()
|
||||
val session = server.connect(collector.send)
|
||||
|
||||
@@ -94,11 +94,11 @@ class ReqResponderServerTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
fun countUsesResponder() =
|
||||
fun countUsesSource() =
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val responder = FixedResponder(listOf(event(1), event(2), event(3)))
|
||||
ReqResponderServer(responder, parentContext = dispatcher).use { server ->
|
||||
val source = FixedSource(listOf(event(1), event(2), event(3)))
|
||||
EventSourceServer(source, parentContext = dispatcher).use { server ->
|
||||
val collector = MessageCollector()
|
||||
val session = server.connect(collector.send)
|
||||
|
||||
@@ -114,17 +114,17 @@ class ReqResponderServerTest {
|
||||
fun approximateCountWithHllReachesTheWire() =
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val responder =
|
||||
object : ReqResponder {
|
||||
override fun respond(filters: List<Filter>): Flow<Event> = flowOf(event(1), event(2))
|
||||
val source =
|
||||
object : EventSource {
|
||||
override fun events(filters: List<Filter>): Flow<Event> = flowOf(event(1), event(2))
|
||||
|
||||
override suspend fun countResult(filters: List<Filter>): CountResult {
|
||||
val hll = HllBuilder(offset = 8)
|
||||
respond(filters).collect { hll.add(it.pubKey) }
|
||||
events(filters).collect { hll.add(it.pubKey) }
|
||||
return hll.toCountResult()
|
||||
}
|
||||
}
|
||||
ReqResponderServer(responder, parentContext = dispatcher).use { server ->
|
||||
EventSourceServer(source, parentContext = dispatcher).use { server ->
|
||||
val collector = MessageCollector()
|
||||
val session = server.connect(collector.send)
|
||||
|
||||
@@ -141,7 +141,7 @@ class ReqResponderServerTest {
|
||||
fun eventPublishIsRejected() =
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
ReqResponderServer(FixedResponder(emptyList()), parentContext = dispatcher).use { server ->
|
||||
EventSourceServer(FixedSource(emptyList()), parentContext = dispatcher).use { server ->
|
||||
val collector = MessageCollector()
|
||||
val session = server.connect(collector.send)
|
||||
|
||||
@@ -155,14 +155,14 @@ class ReqResponderServerTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
fun responderErrorBecomesClosed() =
|
||||
fun sourceErrorBecomesClosed() =
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val responder =
|
||||
object : ReqResponder {
|
||||
override fun respond(filters: List<Filter>): Flow<Event> = flow { throw RuntimeException("backend down") }
|
||||
val source =
|
||||
object : EventSource {
|
||||
override fun events(filters: List<Filter>): Flow<Event> = flow { throw RuntimeException("backend down") }
|
||||
}
|
||||
ReqResponderServer(responder, parentContext = dispatcher).use { server ->
|
||||
EventSourceServer(source, parentContext = dispatcher).use { server ->
|
||||
val collector = MessageCollector()
|
||||
val session = server.connect(collector.send)
|
||||
|
||||
@@ -176,13 +176,13 @@ class ReqResponderServerTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
fun policyGatesTheResponder() =
|
||||
fun policyGatesTheSource() =
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val relay = NormalizedRelayUrl("wss://search.example.com/")
|
||||
val responder = FixedResponder(listOf(event(1)))
|
||||
ReqResponderServer(
|
||||
responder,
|
||||
val source = FixedSource(listOf(event(1)))
|
||||
EventSourceServer(
|
||||
source,
|
||||
policyBuilder = { FullAuthPolicy(relay) },
|
||||
parentContext = dispatcher,
|
||||
).use { server ->
|
||||
@@ -193,7 +193,7 @@ class ReqResponderServerTest {
|
||||
assertTrue(collector.messages[0].contains("\"AUTH\""))
|
||||
assertTrue(OptimizedJsonMapper.fromJsonToMessage(collector.messages[0]) is AuthMessage)
|
||||
|
||||
// A REQ before auth is rejected by the policy, not the responder.
|
||||
// A REQ before auth is rejected by the policy, not the source.
|
||||
session.receive("""["REQ","sub1",{"kinds":[1]}]""")
|
||||
val closed = collector.containing("CLOSED")
|
||||
assertEquals(1, closed.size)
|
||||
+14
-14
@@ -40,9 +40,9 @@ class RelayLimitsServerTest {
|
||||
|
||||
private fun hexId(n: Int): String = n.toString().padStart(64, '0')
|
||||
|
||||
private val emptyResponder =
|
||||
object : ReqResponder {
|
||||
override fun respond(filters: List<Filter>): Flow<Event> = emptyFlow()
|
||||
private val emptySource =
|
||||
object : EventSource {
|
||||
override fun events(filters: List<Filter>): Flow<Event> = emptyFlow()
|
||||
}
|
||||
|
||||
private class Collector {
|
||||
@@ -56,7 +56,7 @@ class RelayLimitsServerTest {
|
||||
fun rejectsOversizedMessageWithNotice() =
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val server = ReqResponderServer(emptyResponder, parentContext = dispatcher, limits = RelayLimits(maxMessageLength = 30))
|
||||
val server = EventSourceServer(emptySource, parentContext = dispatcher, limits = RelayLimits(maxMessageLength = 30))
|
||||
val collector = Collector()
|
||||
val session = server.connect(collector.send)
|
||||
|
||||
@@ -74,12 +74,12 @@ class RelayLimitsServerTest {
|
||||
fun enforcesMaxSubscriptions() =
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
// FixedResponder-style: respond stays open is not needed; emptyFlow EOSEs.
|
||||
val responder =
|
||||
object : ReqResponder {
|
||||
override fun respond(filters: List<Filter>): Flow<Event> = emptyFlow()
|
||||
// An empty flow EOSEs immediately; no need to keep it open.
|
||||
val source =
|
||||
object : EventSource {
|
||||
override fun events(filters: List<Filter>): Flow<Event> = emptyFlow()
|
||||
}
|
||||
val server = ReqResponderServer(responder, parentContext = dispatcher, limits = RelayLimits(maxSubscriptions = 2))
|
||||
val server = EventSourceServer(source, parentContext = dispatcher, limits = RelayLimits(maxSubscriptions = 2))
|
||||
val collector = Collector()
|
||||
val session = server.connect(collector.send)
|
||||
|
||||
@@ -100,20 +100,20 @@ class RelayLimitsServerTest {
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val seen = mutableListOf<Int?>()
|
||||
val responder =
|
||||
object : ReqResponder {
|
||||
override fun respond(filters: List<Filter>): Flow<Event> {
|
||||
val source =
|
||||
object : EventSource {
|
||||
override fun events(filters: List<Filter>): Flow<Event> {
|
||||
seen.add(filters.single().limit)
|
||||
return emptyFlow()
|
||||
}
|
||||
}
|
||||
val server = ReqResponderServer(responder, parentContext = dispatcher, limits = RelayLimits(maxLimit = 100))
|
||||
val server = EventSourceServer(source, parentContext = dispatcher, limits = RelayLimits(maxLimit = 100))
|
||||
val collector = Collector()
|
||||
val session = server.connect(collector.send)
|
||||
|
||||
session.receive("""["REQ","s",{"kinds":[1],"limit":9999}]""")
|
||||
|
||||
assertEquals(listOf<Int?>(100), seen) // the responder saw the clamped limit
|
||||
assertEquals(listOf<Int?>(100), seen) // the source saw the clamped limit
|
||||
server.close()
|
||||
}
|
||||
|
||||
|
||||
+6
-6
@@ -48,9 +48,9 @@ class RelayServerListenerTest {
|
||||
}
|
||||
}
|
||||
|
||||
private val emptyResponder =
|
||||
object : ReqResponder {
|
||||
override fun respond(filters: List<Filter>): Flow<Event> = emptyFlow()
|
||||
private val emptySource =
|
||||
object : EventSource {
|
||||
override fun events(filters: List<Filter>): Flow<Event> = emptyFlow()
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -58,7 +58,7 @@ class RelayServerListenerTest {
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val listener = RecordingListener()
|
||||
val server = ReqResponderServer(emptyResponder, parentContext = dispatcher, listener = listener)
|
||||
val server = EventSourceServer(emptySource, parentContext = dispatcher, listener = listener)
|
||||
|
||||
val a = server.connect {}
|
||||
val b = server.connect {}
|
||||
@@ -83,7 +83,7 @@ class RelayServerListenerTest {
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val listener = RecordingListener()
|
||||
val server = ReqResponderServer(emptyResponder, parentContext = dispatcher, listener = listener)
|
||||
val server = EventSourceServer(emptySource, parentContext = dispatcher, listener = listener)
|
||||
|
||||
val s = server.connect {}
|
||||
s.close()
|
||||
@@ -100,7 +100,7 @@ class RelayServerListenerTest {
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val listener = RecordingListener()
|
||||
val server = ReqResponderServer(emptyResponder, parentContext = dispatcher, listener = listener)
|
||||
val server = EventSourceServer(emptySource, parentContext = dispatcher, listener = listener)
|
||||
|
||||
val s = server.connect {}
|
||||
assertEquals(1L, server.activeConnections)
|
||||
|
||||
Reference in New Issue
Block a user