mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 03:38:23 +00:00
Merge pull request #3133 from vitorpamplona/claude/wizardly-fermat-y7NBf
Move authenticated identity from policy to connection scope
This commit is contained in:
+25
-5
@@ -93,11 +93,13 @@ framing, subscription lifecycle), so there's no hand-written read loop.
|
||||
|
||||
```kotlin
|
||||
class SearchEventSource(private val backend: SearchApi) : EventSource {
|
||||
override fun events(filters: List<Filter>): Flow<Event> = flow {
|
||||
override fun events(ctx: RequestContext, filters: List<Filter>): Flow<Event> = flow {
|
||||
// ctx says who is asking — score/restrict results from their perspective.
|
||||
val viewer = ctx.authenticatedUsers.firstOrNull()
|
||||
filters.forEach { f ->
|
||||
f.search?.let { raw ->
|
||||
val q = SearchQuery.parse(raw)
|
||||
backend.search(q.terms, domain = q.domain, language = q.language)
|
||||
backend.search(q.terms, domain = q.domain, language = q.language, viewer = viewer)
|
||||
.forEach { emit(it) }
|
||||
}
|
||||
}
|
||||
@@ -140,6 +142,24 @@ 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.
|
||||
|
||||
### Caller-aware sources (who is asking)
|
||||
|
||||
Every `events`/`count`/`countResult` call receives a `RequestContext` carrying
|
||||
the connection's `authenticatedUsers` (the pubkeys that completed NIP-42 on this
|
||||
socket) and a stable `connectionId`. That is what makes NIP-42 useful on a
|
||||
non-storage relay: a single shared `EventSource` instance can serve
|
||||
caller-relative results (trust/relevance scored from the viewer's perspective,
|
||||
"for-you" feeds), restricted content (a pubkey's DMs returned only to that
|
||||
pubkey), or per-connection tenancy — without smuggling auth state through a
|
||||
side channel. `ctx.authenticatedUsers` is a live view of the engine-owned
|
||||
connection scope, so a REQ that arrives after the AUTH sees the freshly
|
||||
recorded pubkey(s). (The engine records them on a successful NIP-42 AUTH; the
|
||||
policy reads the same scope to gate.)
|
||||
|
||||
For per-connection state richer than the pubkey (e.g. a backend session token
|
||||
minted in `FullAuthPolicy.authorize`), downcast `ctx.policy` to your own policy
|
||||
subclass and read a typed field — the policy instance is itself per-connection.
|
||||
|
||||
## Policies
|
||||
|
||||
Policies control what clients can do. They validate commands and can rewrite filters.
|
||||
@@ -360,16 +380,16 @@ are omitted), and `CONTENT_TYPE` is `application/nostr+json`.
|
||||
|
||||
## Approximate COUNT (NIP-45 HyperLogLog)
|
||||
|
||||
`COUNT` is answered by `SessionBackend.countResult(filters)` (and
|
||||
`COUNT` is answered by `SessionBackend.countResult(ctx, filters)` (and
|
||||
`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`:
|
||||
|
||||
```kotlin
|
||||
override suspend fun countResult(filters: List<Filter>): CountResult {
|
||||
override suspend fun countResult(ctx: RequestContext, filters: List<Filter>): CountResult {
|
||||
val filter = filters.first()
|
||||
val hll = HyperLogLog.builderFor(filter) ?: return CountResult(count(filters))
|
||||
val hll = HyperLogLog.builderFor(filter) ?: return CountResult(count(ctx, filters))
|
||||
store.query(filter) { event -> hll.add(event.pubKey) }
|
||||
return hll.toCountResult() // count = estimate, approximate = true, hll = registers
|
||||
}
|
||||
|
||||
@@ -0,0 +1,157 @@
|
||||
# Auth identity is connection scope, not policy state
|
||||
|
||||
**Date:** 2026-06-04
|
||||
**Module:** `quartz` — `nip01Core/relay/server`
|
||||
**Status:** Proposed (review before implementing)
|
||||
|
||||
## Problem
|
||||
|
||||
`IRelayPolicy` is a *decision* interface (`accept(...)` → allow/reject/rewrite).
|
||||
But we attached `authenticatedUsers` (the set of pubkeys logged in on a
|
||||
connection) to the policy world via the `AuthScopedPolicy` mixin. That is
|
||||
connection *state*, not a decision. Storing it in the policy forced three
|
||||
artifacts whose only job is to route the state back out as scope:
|
||||
|
||||
- `AuthScopedPolicy` — marker to locate the state,
|
||||
- `PolicyStack.authenticatedUsers` — union to forward it through the composition
|
||||
wrapper (the session's real policy is usually `PolicyStack(LimitsPolicy, …)`),
|
||||
- `RequestContext`'s `as? AuthScopedPolicy` downcast — to read it back.
|
||||
|
||||
`RequestContext` is *already* the connection-scope object (`connectionId`,
|
||||
`authenticatedUsers`). It should **own** the authenticated identity instead of
|
||||
reaching into the policy for it.
|
||||
|
||||
## Why it got merged
|
||||
|
||||
The gating policies need the state to decide: `FullAuthPolicy.accept(ReqCmd)`
|
||||
returns `auth-required` when no one is authenticated. Storing the set in the
|
||||
policy kept `accept()` self-contained. The fix is to let the policy *read* a
|
||||
scope it doesn't *own*.
|
||||
|
||||
## Target model
|
||||
|
||||
- **Connection scope** (engine-owned, per `RelaySession`) owns the mutable
|
||||
`authenticatedUsers` set + `connectionId`. `RequestContext` is its read-only
|
||||
view (what sources already receive).
|
||||
- **Policy stays pure decision.** `FullAuthPolicy` keeps `accept(AuthCmd)`
|
||||
(validate the proof) and gating `accept(ReqCmd/CountCmd/EventCmd)`, but
|
||||
*reads* the scope to gate instead of owning a set.
|
||||
- **The engine performs the single commit.** On a fully-approved AUTH, the
|
||||
engine records the verified pubkey into the scope — one commit point, not a
|
||||
field buried in a policy.
|
||||
|
||||
Deleted by this change: `AuthScopedPolicy`, `PolicyStack.authenticatedUsers`,
|
||||
the `RequestContext` downcast, and `FullAuthPolicy`'s `authenticatedUsers` field
|
||||
/ `isAuthenticated()` / the `.add()` in `onAuthenticated`.
|
||||
|
||||
## The two correctness invariants (must be preserved)
|
||||
|
||||
1. **Only a *verified* identity may enter the scope.** A blind-accept policy
|
||||
(`PassThroughPolicy`/`EmptyPolicy`, which accept `AUTH` without checking a
|
||||
challenge or signature) must record *nothing*. Today this holds because only
|
||||
`FullAuthPolicy.onAuthenticated` commits. We must not regress to "commit
|
||||
whenever `accept` passed."
|
||||
2. **No partial/rolled-back auth.** A thrown `authorize`, or a downstream policy
|
||||
rejecting the AUTH, must leave the connection unauthenticated. Guarded by
|
||||
`authorizeThrowTurnsAuthIntoFailingOk`,
|
||||
`policyRejectingAuthAfterFullAuthLeavesConnectionUnauthenticated`,
|
||||
`failedReAuthKeepsPreviousValidAuthentication`,
|
||||
`commandsRejectedAfterFailedAuthHook`.
|
||||
|
||||
**Hard constraint:** `VerifyPolicy`/`VerifyAuthOnlyPolicy`/`EmptyPolicy` are
|
||||
shared `object` singletons. They must remain stateless — a per-connection scope
|
||||
ref may only be retained by a per-connection policy (built fresh by
|
||||
`policyBuilder`, i.e. `FullAuthPolicy`).
|
||||
|
||||
## Mechanism
|
||||
|
||||
Two small signature changes carry it:
|
||||
|
||||
### A. Read side — inject the scope at connect
|
||||
|
||||
`IRelayPolicy.onConnect(send)` → `onConnect(scope: RequestContext, send)`.
|
||||
|
||||
- `PassThroughPolicy`, `VerifyEventsAndAuthPolicy` (singletons): ignore `scope`
|
||||
(no storage — stays stateless).
|
||||
- `PolicyStack.onConnect`: forward `scope` to each child.
|
||||
- `FullAuthPolicy`: store the read-only `scope` (per-connection, safe) and read
|
||||
`scope.authenticatedUsers` in its gating `accept(...)`.
|
||||
|
||||
### B. Write side — `onAuthenticated` returns the commit decision
|
||||
|
||||
`IRelayPolicy.onAuthenticated(pubKey, event)` → returns `Boolean`
|
||||
(default `false` = "I did not authenticate anyone; record nothing").
|
||||
|
||||
- `FullAuthPolicy`: runs `authorize(...)` (may throw → engine catches → `OK
|
||||
false`), then `return true`. No `.add()`.
|
||||
- `PolicyStack`: run **every** child (side effects) and OR the results, so a
|
||||
chain authenticates iff some verifying member claims it:
|
||||
`policies.fold(false) { acc, p -> p.onAuthenticated(pubKey, event) || acc }`
|
||||
(left operand always evaluated → every hook runs).
|
||||
- Engine (`RelaySession.handleAuth`):
|
||||
```
|
||||
val result = policy.accept(cmd) // proof check across chain
|
||||
if (rejected) { OK false; return }
|
||||
val record = try { policy.onAuthenticated(cmd.event.pubKey, cmd.event) }
|
||||
catch (e) { OK false; return } // authorize threw → nothing recorded
|
||||
if (record) scope.add(cmd.event.pubKey) // single engine-side commit
|
||||
OK true
|
||||
```
|
||||
|
||||
This keeps the write strictly engine-side (invariant 1: only `true`-returning
|
||||
verifying policies cause a commit; PassThrough returns `false`) and the commit
|
||||
after a no-throw `onAuthenticated` (invariant 2).
|
||||
|
||||
## File-by-file
|
||||
|
||||
| File | Change |
|
||||
|------|--------|
|
||||
| `backend/RequestContext.kt` | Drop the `as? AuthScopedPolicy` getter; `authenticatedUsers` becomes a plain backed property. Keep `policy` (for app-state downcast). |
|
||||
| `policies/AuthScopedPolicy.kt` | **Delete.** |
|
||||
| `policies/IRelayPolicy.kt` | `onConnect` gains `scope: RequestContext`; `onAuthenticated` returns `Boolean` (default `false`). |
|
||||
| `policies/FullAuthPolicy.kt` | Remove `authenticatedUsers` field, `isAuthenticated()`, the marker. Store injected `scope`; gate on `scope.authenticatedUsers`; `onAuthenticated` runs `authorize` then `return true`. |
|
||||
| `policies/PolicyStack.kt` | Drop `authenticatedUsers`/marker; forward `scope` in `onConnect`; OR-fold `onAuthenticated`. |
|
||||
| `policies/PassThroughPolicy.kt`, `VerifyPolicy.kt` | `onConnect(scope, send)` — ignore `scope` (stay stateless). |
|
||||
| `server/RelaySession.kt` | Own the mutable `authenticatedUsers` set behind `requestContext`; expose `requestContext` (read view) for observability/tests; pass it to `policy.onConnect`; perform the commit in `handleAuth`. |
|
||||
|
||||
No change to `EventSource`/`SessionBackend`/`EventSourceBackend`/`LiveEventStore`
|
||||
— sources still get `RequestContext` and read `authenticatedUsers` exactly as
|
||||
now; only the *backing* moved.
|
||||
|
||||
## Test mapping (assertions retarget; guarantees unchanged)
|
||||
|
||||
`NostrServerAuthTest` currently inspects `(session.policy as FullAuthPolicy)
|
||||
.authenticatedUsers / .isAuthenticated()`. Retarget to the scope, e.g.
|
||||
`session.requestContext.authenticatedUsers` (expose `requestContext` on
|
||||
`RelaySession`).
|
||||
|
||||
| Test | New path / why it still holds |
|
||||
|------|-------------------------------|
|
||||
| `authSucceedsWithValidEvent` | accept ✓ → `onAuthenticated`→true → engine commits → `requestContext.authenticatedUsers` has pubkey. |
|
||||
| `multipleUsersCanAuthenticate` | two successful AUTHs → engine commits each → set has both. |
|
||||
| `authorizeThrowTurnsAuthIntoFailingOk` | `authorize` throws → engine catches before commit → set empty. |
|
||||
| `policyRejectingAuthAfterFullAuthLeavesConnectionUnauthenticated` | chain `accept` rejected → `onAuthenticated` never called → no commit. |
|
||||
| `failedReAuthKeepsPreviousValidAuthentication` | 2nd accept rejected → no commit → set keeps 1st pubkey. |
|
||||
| `commandsRejectedAfterFailedAuthHook` | unchanged (gating reads empty scope). |
|
||||
| `EventSourceServerTest.sourceSeesAuthenticatedUserInContext` | unchanged — already reads `ctx.authenticatedUsers`. |
|
||||
|
||||
## Alternatives considered
|
||||
|
||||
- **Scope exposes `add()`, verifying policy writes it** (no `onAuthenticated`
|
||||
signature change). Rejected: gives policies write access to scope and splits
|
||||
the commit across N policies; the engine-side single commit is easier to audit
|
||||
against invariant 2.
|
||||
- **Pass `scope` into every `accept(cmd, scope)`** instead of injecting at
|
||||
`onConnect`. Rejected: ~4 methods × ~9 policies of churn for a dependency only
|
||||
`FullAuthPolicy` uses.
|
||||
- **Generic `EventSource<C : RequestContext>`** for a typed `AuthRequestContext`
|
||||
(prior discussion). Out of scope; cascades type params through the shared
|
||||
`SessionBackend`/storage path.
|
||||
|
||||
## Risk
|
||||
|
||||
Low surface, high-sensitivity (signed-in correctness). The four bypass tests are
|
||||
the spec; they pass unchanged in *behavior*, only their assertion target moves.
|
||||
The `onConnect`/`onAuthenticated` signature changes are mechanical but touch
|
||||
every policy + any external `IRelayPolicy` impl (source-breaking — acceptable
|
||||
for this in-tree 2025 SPI, same call as the `EventSource` ctx change).
|
||||
+41
-12
@@ -20,6 +20,7 @@
|
||||
*/
|
||||
package com.vitorpamplona.quartz.nip01Core.relay.server
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.ClosedMessage
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.CountMessage
|
||||
@@ -34,6 +35,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.CloseCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.CountCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.EventCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.ReqCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.backend.RequestContext
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.backend.SessionBackend
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.policies.IRelayPolicy
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.policies.PolicyResult
|
||||
@@ -74,6 +76,27 @@ class RelaySession(
|
||||
) : AutoCloseable {
|
||||
private val subscriptions = LargeCache<String, Job>()
|
||||
|
||||
/**
|
||||
* The authenticated-identity store for this connection. The engine is the
|
||||
* only writer (committed in [handleAuth] on a successful NIP-42 AUTH); the
|
||||
* policy and the data plane read it through [requestContext].
|
||||
*/
|
||||
private val authenticatedUsers = mutableSetOf<HexKey>()
|
||||
|
||||
/**
|
||||
* The per-connection scope. Handed to the [policy] at connect (so gating
|
||||
* policies can read the authenticated users) and to the [store] on every
|
||||
* REQ/COUNT (so a source can see who is asking). [RequestContext.authenticatedUsers]
|
||||
* is a live view of [authenticatedUsers], so a REQ after a NIP-42 AUTH sees
|
||||
* the freshly recorded pubkey(s).
|
||||
*/
|
||||
val requestContext: RequestContext =
|
||||
object : RequestContext {
|
||||
override val connectionId = id
|
||||
override val policy = this@RelaySession.policy
|
||||
override val authenticatedUsers: Set<HexKey> get() = this@RelaySession.authenticatedUsers
|
||||
}
|
||||
|
||||
/** NIP-77 negentropy state for this connection. */
|
||||
private val negentropy = NegSessionRegistry(store, ::send, negentropySettings)
|
||||
|
||||
@@ -189,7 +212,7 @@ class RelaySession(
|
||||
|
||||
val countResult =
|
||||
try {
|
||||
store.countResult(filters)
|
||||
store.countResult(requestContext, filters)
|
||||
} catch (e: CancellationException) {
|
||||
throw e
|
||||
} catch (e: Exception) {
|
||||
@@ -210,16 +233,21 @@ class RelaySession(
|
||||
|
||||
// The whole policy chain validated the AUTH. onAuthenticated runs any
|
||||
// post-verification I/O (e.g. exchanging the verified event for a
|
||||
// backend token) AND is where a policy commits the authentication, so a
|
||||
// throw here cleanly fails the login — nothing was committed to undo.
|
||||
try {
|
||||
policy.onAuthenticated(cmd.event.pubKey, cmd.event)
|
||||
} catch (e: CancellationException) {
|
||||
throw e
|
||||
} catch (e: Exception) {
|
||||
send(OkMessage.rejected(cmd.event.id, MachineReadablePrefix.ERROR, e.message ?: "authentication failed"))
|
||||
return
|
||||
}
|
||||
// backend token) and votes on whether to record the identity. A throw
|
||||
// here cleanly fails the login — the engine records nothing.
|
||||
val record =
|
||||
try {
|
||||
policy.onAuthenticated(cmd.event)
|
||||
} catch (e: CancellationException) {
|
||||
throw e
|
||||
} catch (e: Exception) {
|
||||
send(OkMessage.rejected(cmd.event.id, MachineReadablePrefix.ERROR, e.message ?: "authentication failed"))
|
||||
return
|
||||
}
|
||||
|
||||
// Single, engine-side commit into the connection scope — after the full
|
||||
// chain approved and a verifying policy voted to record.
|
||||
if (record) authenticatedUsers.add(cmd.event.pubKey)
|
||||
|
||||
send(OkMessage(cmd.event.id, true, ""))
|
||||
}
|
||||
@@ -254,6 +282,7 @@ class RelaySession(
|
||||
scope.launch {
|
||||
try {
|
||||
store.query(
|
||||
ctx = requestContext,
|
||||
filters = filters,
|
||||
onEach = { event ->
|
||||
if (policy.canSendToSession(event)) {
|
||||
@@ -285,7 +314,7 @@ class RelaySession(
|
||||
}
|
||||
|
||||
init {
|
||||
policy.onConnect(::send)
|
||||
policy.onConnect(requestContext, ::send)
|
||||
}
|
||||
|
||||
companion object {
|
||||
|
||||
+31
-11
@@ -37,17 +37,27 @@ import kotlinx.coroutines.flow.count
|
||||
*
|
||||
* ```
|
||||
* class SearchEventSource(private val backend: SearchApi) : EventSource {
|
||||
* override fun events(filters: List<Filter>): Flow<Event> = flow {
|
||||
* override fun events(ctx: RequestContext, filters: List<Filter>): Flow<Event> = flow {
|
||||
* // Tailor the answer to the caller: NIP-42 already told us who they are.
|
||||
* val viewer = ctx.authenticatedUsers.firstOrNull()
|
||||
* for (f in filters) {
|
||||
* f.search?.let { raw ->
|
||||
* val query = SearchQuery.parse(raw)
|
||||
* backend.search(query.terms, query.language).forEach { emit(it) }
|
||||
* backend.search(query.terms, query.language, viewer).forEach { emit(it) }
|
||||
* }
|
||||
* }
|
||||
* }
|
||||
* }
|
||||
* ```
|
||||
*
|
||||
* ## Who is asking
|
||||
*
|
||||
* Every call receives a [RequestContext] carrying the connection's
|
||||
* [RequestContext.authenticatedUsers] (and [RequestContext.connectionId]). A
|
||||
* single shared [EventSource] instance can therefore serve caller-relative
|
||||
* results, restricted content, or per-connection tenancy without smuggling auth
|
||||
* state in through a side channel — see [RequestContext].
|
||||
*
|
||||
* ## EOSE semantics
|
||||
*
|
||||
* The engine sends `EOSE` when the returned [Flow] **completes**. A source
|
||||
@@ -62,18 +72,25 @@ import kotlinx.coroutines.flow.count
|
||||
*/
|
||||
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.
|
||||
* Returns the events matching [filters] for the caller described by [ctx].
|
||||
* 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; [ctx] carries the connection's authenticated identity.
|
||||
*/
|
||||
fun events(filters: List<Filter>): Flow<Event>
|
||||
fun events(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
): Flow<Event>
|
||||
|
||||
/**
|
||||
* Answers a NIP-45 COUNT. The default counts the events produced by
|
||||
* [events]; override it when the backend can count without materializing
|
||||
* every event.
|
||||
* Answers a NIP-45 COUNT for the caller described by [ctx]. The default
|
||||
* counts the events produced by [events]; override it when the backend can
|
||||
* count without materializing every event.
|
||||
*/
|
||||
suspend fun count(filters: List<Filter>): Int = events(filters).count()
|
||||
suspend fun count(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
): Int = events(ctx, filters).count()
|
||||
|
||||
/**
|
||||
* Answers a NIP-45 COUNT, optionally approximate and/or carrying a
|
||||
@@ -81,5 +98,8 @@ interface EventSource {
|
||||
* override to return `approximate`/`hll` (see
|
||||
* [com.vitorpamplona.quartz.nip45Count.HllBuilder]).
|
||||
*/
|
||||
suspend fun countResult(filters: List<Filter>): CountResult = CountResult(count(filters))
|
||||
suspend fun countResult(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
): CountResult = CountResult(count(ctx, filters))
|
||||
}
|
||||
|
||||
+10
-3
@@ -36,15 +36,22 @@ class EventSourceBackend(
|
||||
private val source: EventSource,
|
||||
) : SessionBackend {
|
||||
override suspend fun query(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
onEach: (Event) -> Unit,
|
||||
onEose: () -> Unit,
|
||||
) {
|
||||
source.events(filters).collect { onEach(it) }
|
||||
source.events(ctx, filters).collect { onEach(it) }
|
||||
onEose()
|
||||
}
|
||||
|
||||
override suspend fun count(filters: List<Filter>): Int = source.count(filters)
|
||||
override suspend fun count(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
): Int = source.count(ctx, filters)
|
||||
|
||||
override suspend fun countResult(filters: List<Filter>): CountResult = source.countResult(filters)
|
||||
override suspend fun countResult(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
): CountResult = source.countResult(ctx, filters)
|
||||
}
|
||||
|
||||
+5
-1
@@ -126,6 +126,7 @@ class LiveEventStore(
|
||||
}
|
||||
|
||||
override suspend fun query(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
onEach: (Event) -> Unit,
|
||||
onEose: () -> Unit,
|
||||
@@ -191,7 +192,10 @@ class LiveEventStore(
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun count(filters: List<Filter>): Int = store.count(filters)
|
||||
override suspend fun count(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
): Int = store.count(filters)
|
||||
|
||||
/**
|
||||
* One-shot snapshot query. Used by NIP-77 negentropy: the server
|
||||
|
||||
+69
@@ -0,0 +1,69 @@
|
||||
/*
|
||||
* 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.backend
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.policies.IRelayPolicy
|
||||
|
||||
/**
|
||||
* The per-connection context handed to an [EventSource] (and a [SessionBackend])
|
||||
* on every REQ/COUNT. It tells the data plane *who* is asking, so a single
|
||||
* shared source can tailor its answer to the caller without smuggling state in
|
||||
* through a side channel.
|
||||
*
|
||||
* This is the connection scope: the engine owns it (one per
|
||||
* [com.vitorpamplona.quartz.nip01Core.relay.server.RelaySession]) and records
|
||||
* the authenticated pubkey(s) into it on a successful NIP-42 AUTH. That is what
|
||||
* makes NIP-42 useful on a non-storage relay — [RequestContext] is the path
|
||||
* that carries the caller's identity to the code that produces the events. With
|
||||
* it, the same `EventSource` instance can serve:
|
||||
*
|
||||
* - **caller-relative results** — trust/relevance scored from
|
||||
* [authenticatedUsers]'s perspective ("for-you" feeds, follow-aware search);
|
||||
* - **restricted content** — a pubkey's DMs/private events returned only to
|
||||
* that authenticated pubkey;
|
||||
* - **paid / allow-listed** sets that depend on the authenticated identity;
|
||||
* - **per-connection** quotas/tenancy keyed by [connectionId].
|
||||
*
|
||||
* For per-connection application state beyond the pubkey (e.g. a backend session
|
||||
* token minted in [com.vitorpamplona.quartz.nip01Core.relay.server.policies.FullAuthPolicy.authorize]),
|
||||
* downcast [policy] to your own [IRelayPolicy] subclass and read a typed field —
|
||||
* the policy instance is itself per-connection, so it is the natural typed bag.
|
||||
*/
|
||||
interface RequestContext {
|
||||
/**
|
||||
* Stable, process-unique id of the connection this request arrived on
|
||||
* (the owning [com.vitorpamplona.quartz.nip01Core.relay.server.RelaySession.id]).
|
||||
* Use it to key per-connection state a shared source keeps on the side.
|
||||
*/
|
||||
val connectionId: Long
|
||||
|
||||
/** The connection's policy, for typed access to per-connection state. */
|
||||
val policy: IRelayPolicy
|
||||
|
||||
/**
|
||||
* The pubkeys that have authenticated on this connection via NIP-42. Empty
|
||||
* when the connection is unauthenticated. Backed by the engine-owned scope
|
||||
* and read live, so a REQ that arrives after a successful AUTH sees the
|
||||
* freshly recorded pubkey(s).
|
||||
*/
|
||||
val authenticatedUsers: Set<HexKey>
|
||||
}
|
||||
+15
-7
@@ -45,19 +45,24 @@ import com.vitorpamplona.quartz.nip01Core.store.IdAndTime
|
||||
*/
|
||||
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 source returns after [onEose].
|
||||
* Answers a REQ on behalf of the caller described by [ctx]. 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 source returns after
|
||||
* [onEose].
|
||||
*/
|
||||
suspend fun query(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
onEach: (Event) -> Unit,
|
||||
onEose: () -> Unit,
|
||||
)
|
||||
|
||||
/** Answers a NIP-45 COUNT with an exact cardinality. */
|
||||
suspend fun count(filters: List<Filter>): Int
|
||||
/** Answers a NIP-45 COUNT with an exact cardinality for the caller in [ctx]. */
|
||||
suspend fun count(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
): Int
|
||||
|
||||
/**
|
||||
* Answers a NIP-45 COUNT, allowing an approximate result and/or a
|
||||
@@ -65,7 +70,10 @@ interface SessionBackend {
|
||||
* The default returns the exact [count] with `approximate = false`; override
|
||||
* to return `approximate`/`hll`.
|
||||
*/
|
||||
suspend fun countResult(filters: List<Filter>): CountResult = CountResult(count(filters))
|
||||
suspend fun countResult(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
): CountResult = CountResult(count(ctx, filters))
|
||||
|
||||
/**
|
||||
* Handles an EVENT publish, reporting the per-event outcome through
|
||||
|
||||
+40
-27
@@ -29,6 +29,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.CountCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.EventCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.ReqCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.backend.RequestContext
|
||||
import com.vitorpamplona.quartz.nip40Expiration.isExpired
|
||||
import com.vitorpamplona.quartz.nip42RelayAuth.RelayAuthEvent
|
||||
import com.vitorpamplona.quartz.utils.RandomInstance
|
||||
@@ -40,11 +41,13 @@ import com.vitorpamplona.quartz.utils.TimeUtils
|
||||
*
|
||||
* Implements the full NIP-42 challenge/verify handshake: [onConnect] sends the
|
||||
* [challenge] and [accept] (AuthCmd) validates the returned event (expiration,
|
||||
* freshness, challenge match, relay match). Crucially, [accept] does NOT mutate
|
||||
* state — the pubkey is recorded in [authenticatedUsers] only by [onAuthenticated],
|
||||
* which the engine calls once [accept] *and* the rest of the policy chain have
|
||||
* approved the AUTH. That single, late commit is why there is no rollback to
|
||||
* reason about: a rejected AUTH simply never reaches it.
|
||||
* freshness, challenge match, relay match). This policy runs the auth *logic*
|
||||
* but does not *own* the authenticated-identity store: the engine-owned
|
||||
* connection [scope] holds it. [accept] does NOT mutate state, and
|
||||
* [onAuthenticated] only votes `true` (after [authorize]) — the engine performs
|
||||
* the single recording into the scope once the whole chain has approved. Gating
|
||||
* decisions read [scope].`authenticatedUsers`. A rejected AUTH simply never
|
||||
* reaches [onAuthenticated], so there is no rollback to reason about.
|
||||
*
|
||||
* To bridge to an external auth system, override [authorize] (a `suspend` hook)
|
||||
* and do the post-verification I/O there — e.g. exchange the verified event for
|
||||
@@ -57,13 +60,30 @@ open class FullAuthPolicy(
|
||||
/** The challenge string sent to this client for NIP-42 authentication. */
|
||||
val challenge: String = RandomInstance.randomChars(32)
|
||||
|
||||
/** Set of pubkeys that have successfully authenticated on this session. */
|
||||
val authenticatedUsers = mutableSetOf<HexKey>()
|
||||
/**
|
||||
* The engine-owned connection scope, captured at [onConnect]. Read-only
|
||||
* here: this policy reads [RequestContext.authenticatedUsers] to gate, while
|
||||
* the engine is the only writer. Held safely because a [FullAuthPolicy] is
|
||||
* built fresh per connection.
|
||||
*/
|
||||
private lateinit var scope: RequestContext
|
||||
|
||||
/** Returns true if at least one pubkey has authenticated. */
|
||||
/**
|
||||
* The pubkeys authenticated on this connection, read from the engine-owned
|
||||
* scope. Exposed to subclasses so they can gate or rewrite on the caller's
|
||||
* identity (restricted content, caller-relative filters) — the same set the
|
||||
* data plane sees via [RequestContext.authenticatedUsers].
|
||||
*/
|
||||
protected val authenticatedUsers: Set<HexKey> get() = scope.authenticatedUsers
|
||||
|
||||
/** Returns true if at least one pubkey has authenticated on this connection. */
|
||||
fun isAuthenticated(): Boolean = authenticatedUsers.isNotEmpty()
|
||||
|
||||
override fun onConnect(send: (Message) -> Unit) {
|
||||
override fun onConnect(
|
||||
scope: RequestContext,
|
||||
send: (Message) -> Unit,
|
||||
) {
|
||||
this.scope = scope
|
||||
send(AuthMessage(challenge))
|
||||
}
|
||||
|
||||
@@ -90,31 +110,24 @@ open class FullAuthPolicy(
|
||||
}
|
||||
|
||||
/**
|
||||
* Commits the authentication. The engine calls this only after [accept] and
|
||||
* the whole policy chain have approved the AUTH, so this is the single point
|
||||
* where the pubkey is recorded. It runs [authorize] first (which may throw
|
||||
* to reject) and records the pubkey only on success — the connection is
|
||||
* never left authenticated behind a failing `OK`. `final`: override
|
||||
* [authorize], not this.
|
||||
* Votes to record the authentication. The engine calls this only after
|
||||
* [accept] and the whole policy chain have approved the AUTH. It runs
|
||||
* [authorize] first (which may throw to reject — the engine then records
|
||||
* nothing) and returns `true` so the engine records `event.pubKey` into the
|
||||
* connection scope. `final`: override [authorize], not this.
|
||||
*/
|
||||
final override suspend fun onAuthenticated(
|
||||
pubKey: HexKey,
|
||||
event: RelayAuthEvent,
|
||||
) {
|
||||
authorize(pubKey, event)
|
||||
authenticatedUsers.add(pubKey)
|
||||
final override suspend fun onAuthenticated(event: RelayAuthEvent): Boolean {
|
||||
authorize(event)
|
||||
return true
|
||||
}
|
||||
|
||||
/**
|
||||
* Hook for external authorization once the NIP-42 proof checks out — e.g.
|
||||
* exchange [event] for a backend session token. Throw to reject the login
|
||||
* (the AUTH becomes `OK false` and the pubkey is not recorded). Runs before
|
||||
* the pubkey is committed. The default does nothing.
|
||||
* (the AUTH becomes `OK false` and `event.pubKey` is not recorded). Runs
|
||||
* before the pubkey is committed. The default does nothing.
|
||||
*/
|
||||
open suspend fun authorize(
|
||||
pubKey: HexKey,
|
||||
event: RelayAuthEvent,
|
||||
) {}
|
||||
open suspend fun authorize(event: RelayAuthEvent) {}
|
||||
|
||||
override fun accept(cmd: EventCmd): PolicyResult<EventCmd> =
|
||||
if (isAuthenticated()) {
|
||||
|
||||
+28
-16
@@ -21,20 +21,30 @@
|
||||
package com.vitorpamplona.quartz.nip01Core.relay.server.policies
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.AuthCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.Command
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.CountCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.EventCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.ReqCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.backend.RequestContext
|
||||
import com.vitorpamplona.quartz.nip42RelayAuth.RelayAuthEvent
|
||||
|
||||
/**
|
||||
* Defines custom behavior for this relay.
|
||||
*/
|
||||
interface IRelayPolicy {
|
||||
fun onConnect(send: (Message) -> Unit)
|
||||
/**
|
||||
* Called once when the connection opens. [scope] is the engine-owned,
|
||||
* read-only connection scope (id + authenticated users) — a per-connection
|
||||
* policy may retain it to make later auth-aware decisions; shared singleton
|
||||
* policies must ignore it and stay stateless. [send] pushes a message to the
|
||||
* client (e.g. a NIP-42 AUTH challenge).
|
||||
*/
|
||||
fun onConnect(
|
||||
scope: RequestContext,
|
||||
send: (Message) -> Unit,
|
||||
)
|
||||
|
||||
/**
|
||||
* Evaluates whether an incoming EVENT command should be accepted.
|
||||
@@ -70,24 +80,26 @@ interface IRelayPolicy {
|
||||
|
||||
/**
|
||||
* Called once an AUTH command has been [accept]ed by this policy *and* the
|
||||
* rest of the policy chain, before the success `OK` is sent. This is where
|
||||
* a policy commits the authentication and/or runs post-verification side
|
||||
* effects that need network or disk I/O — e.g. exchanging the verified
|
||||
* NIP-42 event for a backend session token — without leaking that logic into
|
||||
* the transport layer.
|
||||
* rest of the policy chain, before the success `OK` is sent. Run any
|
||||
* post-verification side effects that need network or disk I/O here — e.g.
|
||||
* exchanging the verified NIP-42 event for a backend session token — without
|
||||
* leaking that logic into the transport layer.
|
||||
*
|
||||
* The engine — not the policy — owns the authenticated-identity store. The
|
||||
* return value is this policy's vote on whether `event.pubKey` should be
|
||||
* recorded as authenticated on the connection: return `true` only if this
|
||||
* policy actually verified the identity. The default returns `false`, so a
|
||||
* policy that does not authenticate (e.g. a pass-through or a blind-accept)
|
||||
* never causes an unverified pubkey to be recorded.
|
||||
*
|
||||
* Because it runs only after the whole chain approved the AUTH, throwing
|
||||
* here cleanly fails the login: the AUTH becomes `OK false` and, since the
|
||||
* commit lives here too, the connection is never left authenticated. The
|
||||
* default implementation does nothing.
|
||||
* here cleanly fails the login: the AUTH becomes `OK false` and the engine
|
||||
* records nothing.
|
||||
*
|
||||
* @param pubKey The pubkey being authenticated.
|
||||
* @param event The verified NIP-42 auth event.
|
||||
* @param event The verified NIP-42 auth event (its signer is the identity).
|
||||
* @return `true` to have the engine record `event.pubKey` as authenticated.
|
||||
*/
|
||||
suspend fun onAuthenticated(
|
||||
pubKey: HexKey,
|
||||
event: RelayAuthEvent,
|
||||
) {}
|
||||
suspend fun onAuthenticated(event: RelayAuthEvent): Boolean = false
|
||||
|
||||
/**
|
||||
* Inspects a raw inbound message before it is parsed. Return a reason
|
||||
|
||||
+5
-1
@@ -26,6 +26,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.AuthCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.CountCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.EventCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.ReqCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.backend.RequestContext
|
||||
|
||||
/**
|
||||
* Convenience base that accepts everything by default. Subclasses
|
||||
@@ -37,7 +38,10 @@ import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.ReqCmd
|
||||
* instantiate this directly.
|
||||
*/
|
||||
open class PassThroughPolicy : IRelayPolicy {
|
||||
override fun onConnect(send: (Message) -> Unit) {}
|
||||
override fun onConnect(
|
||||
scope: RequestContext,
|
||||
send: (Message) -> Unit,
|
||||
) {}
|
||||
|
||||
override fun accept(cmd: EventCmd): PolicyResult<EventCmd> = PolicyResult.Accepted(cmd)
|
||||
|
||||
|
||||
+10
-8
@@ -21,13 +21,13 @@
|
||||
package com.vitorpamplona.quartz.nip01Core.relay.server.policies
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.AuthCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.Command
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.CountCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.EventCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.ReqCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.backend.RequestContext
|
||||
import com.vitorpamplona.quartz.nip42RelayAuth.RelayAuthEvent
|
||||
|
||||
class PolicyStack(
|
||||
@@ -35,8 +35,11 @@ class PolicyStack(
|
||||
) : IRelayPolicy {
|
||||
val policies = policies.toList()
|
||||
|
||||
override fun onConnect(send: (Message) -> Unit) {
|
||||
policies.forEach { it.onConnect(send) }
|
||||
override fun onConnect(
|
||||
scope: RequestContext,
|
||||
send: (Message) -> Unit,
|
||||
) {
|
||||
policies.forEach { it.onConnect(scope, send) }
|
||||
}
|
||||
|
||||
override fun accept(cmd: EventCmd) = runPolicies(cmd) { p, c -> p.accept(c) }
|
||||
@@ -47,11 +50,10 @@ class PolicyStack(
|
||||
|
||||
override fun accept(cmd: AuthCmd) = runPolicies(cmd) { p, c -> p.accept(c) }
|
||||
|
||||
override suspend fun onAuthenticated(
|
||||
pubKey: HexKey,
|
||||
event: RelayAuthEvent,
|
||||
) {
|
||||
policies.forEach { it.onAuthenticated(pubKey, event) }
|
||||
override suspend fun onAuthenticated(event: RelayAuthEvent): Boolean {
|
||||
// Run every member (side effects) and record iff any one verified the
|
||||
// identity. `fold` keeps the call on the left so no member is skipped.
|
||||
return policies.fold(false) { recorded, p -> p.onAuthenticated(event) || recorded }
|
||||
}
|
||||
|
||||
override fun acceptMessage(message: String): String? = policies.firstNotNullOfOrNull { it.acceptMessage(message) }
|
||||
|
||||
+5
-1
@@ -27,6 +27,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.AuthCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.CountCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.EventCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.ReqCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.backend.RequestContext
|
||||
|
||||
/**
|
||||
* Verifies the Schnorr signature + id hash of every incoming
|
||||
@@ -43,7 +44,10 @@ import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.ReqCmd
|
||||
open class VerifyEventsAndAuthPolicy(
|
||||
private val verifyEvents: Boolean,
|
||||
) : IRelayPolicy {
|
||||
override fun onConnect(send: (Message) -> Unit) { }
|
||||
override fun onConnect(
|
||||
scope: RequestContext,
|
||||
send: (Message) -> Unit,
|
||||
) { }
|
||||
|
||||
override fun accept(cmd: EventCmd) =
|
||||
if (!verifyEvents || cmd.event.verify()) {
|
||||
|
||||
+80
-5
@@ -26,11 +26,15 @@ import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.AuthMessage
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.CountResult
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.EoseMessage
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.EventMessage
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.AuthCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.backend.EventSource
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.backend.RequestContext
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.policies.FullAuthPolicy
|
||||
import com.vitorpamplona.quartz.nip42RelayAuth.RelayAuthEvent
|
||||
import com.vitorpamplona.quartz.nip45Count.HllBuilder
|
||||
import com.vitorpamplona.quartz.utils.TimeUtils
|
||||
import kotlinx.coroutines.ExperimentalCoroutinesApi
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.collect
|
||||
@@ -70,7 +74,10 @@ class EventSourceServerTest {
|
||||
private class FixedSource(
|
||||
private val events: List<Event>,
|
||||
) : EventSource {
|
||||
override fun events(filters: List<Filter>): Flow<Event> = flowOf(*events.toTypedArray())
|
||||
override fun events(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
): Flow<Event> = flowOf(*events.toTypedArray())
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -117,11 +124,17 @@ class EventSourceServerTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val source =
|
||||
object : EventSource {
|
||||
override fun events(filters: List<Filter>): Flow<Event> = flowOf(event(1), event(2))
|
||||
override fun events(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
): Flow<Event> = flowOf(event(1), event(2))
|
||||
|
||||
override suspend fun countResult(filters: List<Filter>): CountResult {
|
||||
override suspend fun countResult(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
): CountResult {
|
||||
val hll = HllBuilder(offset = 8)
|
||||
events(filters).collect { hll.add(it.pubKey) }
|
||||
events(ctx, filters).collect { hll.add(it.pubKey) }
|
||||
return hll.toCountResult()
|
||||
}
|
||||
}
|
||||
@@ -161,7 +174,10 @@ class EventSourceServerTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val source =
|
||||
object : EventSource {
|
||||
override fun events(filters: List<Filter>): Flow<Event> = flow { throw RuntimeException("backend down") }
|
||||
override fun events(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
): Flow<Event> = flow { throw RuntimeException("backend down") }
|
||||
}
|
||||
EventSourceServer(source, parentContext = dispatcher).use { server ->
|
||||
val collector = MessageCollector()
|
||||
@@ -201,4 +217,63 @@ class EventSourceServerTest {
|
||||
assertTrue(closed[0].contains("auth-required:"))
|
||||
}
|
||||
}
|
||||
|
||||
/** Builds a kind 22242 auth event; FullAuthPolicy checks challenge/relay/time, not the signature. */
|
||||
private fun authEvent(
|
||||
challenge: String,
|
||||
relay: String,
|
||||
) = RelayAuthEvent(
|
||||
id = hexId(99),
|
||||
pubKey = pubkey,
|
||||
createdAt = TimeUtils.now(),
|
||||
tags =
|
||||
arrayOf(
|
||||
arrayOf("relay", relay),
|
||||
arrayOf("challenge", challenge),
|
||||
),
|
||||
content = "",
|
||||
sig = sig,
|
||||
)
|
||||
|
||||
private fun authJson(event: RelayAuthEvent) = OptimizedJsonMapper.toJson(AuthCmd(event))
|
||||
|
||||
@Test
|
||||
fun sourceSeesAuthenticatedUserInContext() =
|
||||
runTest {
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val relay = NormalizedRelayUrl("wss://search.example.com/")
|
||||
|
||||
// A caller-aware source: it records who the engine says is asking.
|
||||
var seenViewers: Set<String>? = null
|
||||
val source =
|
||||
object : EventSource {
|
||||
override fun events(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
): Flow<Event> {
|
||||
seenViewers = ctx.authenticatedUsers
|
||||
return flowOf(event(1))
|
||||
}
|
||||
}
|
||||
|
||||
EventSourceServer(
|
||||
source,
|
||||
policyBuilder = { FullAuthPolicy(relay) },
|
||||
parentContext = dispatcher,
|
||||
).use { server ->
|
||||
val collector = MessageCollector()
|
||||
val session = server.connect(collector.send)
|
||||
|
||||
// Complete the NIP-42 handshake using the engine's challenge.
|
||||
val challenge = (OptimizedJsonMapper.fromJsonToMessage(collector.messages[0]) as AuthMessage).challenge
|
||||
session.receive(authJson(authEvent(challenge, relay.url)))
|
||||
|
||||
// A REQ after auth must hand the authenticated pubkey to the source.
|
||||
session.receive("""["REQ","sub1",{"kinds":[1]}]""")
|
||||
|
||||
assertEquals(setOf(pubkey), seenViewers)
|
||||
val events = collector.parsed().filterIsInstance<EventMessage>()
|
||||
assertEquals(1, events.size)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+38
-19
@@ -168,7 +168,7 @@ class NostrServerAuthTest {
|
||||
assertEquals(1, okMessages.size)
|
||||
assertTrue(okMessages[0].contains(",true,"))
|
||||
assertTrue((session.policy as FullAuthPolicy).isAuthenticated())
|
||||
assertTrue(session.policy.authenticatedUsers.contains(pubkey))
|
||||
assertTrue(session.requestContext.authenticatedUsers.contains(pubkey))
|
||||
|
||||
server.close()
|
||||
}
|
||||
@@ -313,7 +313,7 @@ class NostrServerAuthTest {
|
||||
assertTrue(okMessages[0].contains(",true,"))
|
||||
assertTrue(okMessages[1].contains(",true,"))
|
||||
|
||||
val authedPubkeys = (session.policy as FullAuthPolicy).authenticatedUsers
|
||||
val authedPubkeys = session.requestContext.authenticatedUsers
|
||||
assertEquals(2, authedPubkeys.size)
|
||||
assertTrue(authedPubkeys.contains(pubkey))
|
||||
assertTrue(authedPubkeys.contains(pubkey2))
|
||||
@@ -321,6 +321,34 @@ class NostrServerAuthTest {
|
||||
server.close()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun authenticationIsScopedPerConnection() =
|
||||
runTest {
|
||||
// Two connections on the SAME server. Each must see only the pubkey
|
||||
// that authenticated on it — never the union across the relay.
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
val server = createServer(dispatcher = dispatcher)
|
||||
|
||||
val c1 = MessageCollector()
|
||||
val s1 = server.connect(c1.sendCallback)
|
||||
val ch1 = (OptimizedJsonMapper.fromJsonToMessage(c1.messages[0]) as AuthMessage).challenge
|
||||
s1.receive(authJson(authEvent(challenge = ch1, pubKey = pubkey)))
|
||||
|
||||
val c2 = MessageCollector()
|
||||
val s2 = server.connect(c2.sendCallback)
|
||||
val ch2 = (OptimizedJsonMapper.fromJsonToMessage(c2.messages[0]) as AuthMessage).challenge
|
||||
s2.receive(authJson(authEvent(challenge = ch2, pubKey = pubkey2)))
|
||||
|
||||
// Each connection's scope holds exactly its own authenticated user.
|
||||
assertEquals(setOf(pubkey), s1.requestContext.authenticatedUsers)
|
||||
assertEquals(setOf(pubkey2), s2.requestContext.authenticatedUsers)
|
||||
// Cross-check: neither leaks the other's identity.
|
||||
assertFalse(s1.requestContext.authenticatedUsers.contains(pubkey2))
|
||||
assertFalse(s2.requestContext.authenticatedUsers.contains(pubkey))
|
||||
|
||||
server.close()
|
||||
}
|
||||
|
||||
// -- NIP-42: requireAuth ---------------------------------------------------
|
||||
|
||||
@Test
|
||||
@@ -545,11 +573,8 @@ class NostrServerAuthTest {
|
||||
var hookPubkey: String? = null
|
||||
val policy =
|
||||
object : FullAuthPolicy(relayUrl) {
|
||||
override suspend fun authorize(
|
||||
pubKey: String,
|
||||
event: RelayAuthEvent,
|
||||
) {
|
||||
hookPubkey = pubKey
|
||||
override suspend fun authorize(event: RelayAuthEvent) {
|
||||
hookPubkey = event.pubKey
|
||||
}
|
||||
}
|
||||
|
||||
@@ -574,10 +599,7 @@ class NostrServerAuthTest {
|
||||
runTest {
|
||||
val policy =
|
||||
object : FullAuthPolicy(relayUrl) {
|
||||
override suspend fun authorize(
|
||||
pubKey: String,
|
||||
event: RelayAuthEvent,
|
||||
): Unit = throw IllegalStateException("backend rejected user")
|
||||
override suspend fun authorize(event: RelayAuthEvent): Unit = throw IllegalStateException("backend rejected user")
|
||||
}
|
||||
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
@@ -596,7 +618,7 @@ class NostrServerAuthTest {
|
||||
// still-authenticated connection would be an auth bypass.
|
||||
val authPolicy = session.policy as FullAuthPolicy
|
||||
assertFalse(authPolicy.isAuthenticated())
|
||||
assertFalse(authPolicy.authenticatedUsers.contains(pubkey))
|
||||
assertFalse(session.requestContext.authenticatedUsers.contains(pubkey))
|
||||
|
||||
server.close()
|
||||
}
|
||||
@@ -606,10 +628,7 @@ class NostrServerAuthTest {
|
||||
runTest {
|
||||
val policy =
|
||||
object : FullAuthPolicy(relayUrl) {
|
||||
override suspend fun authorize(
|
||||
pubKey: String,
|
||||
event: RelayAuthEvent,
|
||||
): Unit = throw IllegalStateException("backend rejected user")
|
||||
override suspend fun authorize(event: RelayAuthEvent): Unit = throw IllegalStateException("backend rejected user")
|
||||
}
|
||||
|
||||
val dispatcher = UnconfinedTestDispatcher(testScheduler)
|
||||
@@ -653,7 +672,7 @@ class NostrServerAuthTest {
|
||||
assertEquals(1, ok.size)
|
||||
assertTrue(ok[0].contains(",false,"))
|
||||
assertFalse(auth.isAuthenticated())
|
||||
assertFalse(auth.authenticatedUsers.contains(pubkey))
|
||||
assertFalse(session.requestContext.authenticatedUsers.contains(pubkey))
|
||||
|
||||
// And a privileged REQ is still gated.
|
||||
session.receive("""["REQ","sub1",{"kinds":[1]}]""")
|
||||
@@ -686,12 +705,12 @@ class NostrServerAuthTest {
|
||||
val msg = OptimizedJsonMapper.fromJsonToMessage(collector.messages[0]) as AuthMessage
|
||||
|
||||
session.receive(authJson(authEvent(challenge = msg.challenge)))
|
||||
assertTrue(auth.authenticatedUsers.contains(pubkey))
|
||||
assertTrue(session.requestContext.authenticatedUsers.contains(pubkey))
|
||||
|
||||
// Second AUTH for the SAME pubkey, rejected downstream.
|
||||
session.receive(authJson(authEvent(challenge = msg.challenge)))
|
||||
assertTrue(auth.isAuthenticated())
|
||||
assertTrue(auth.authenticatedUsers.contains(pubkey))
|
||||
assertTrue(session.requestContext.authenticatedUsers.contains(pubkey))
|
||||
|
||||
server.close()
|
||||
}
|
||||
|
||||
+13
-3
@@ -23,6 +23,7 @@ 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.backend.EventSource
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.backend.RequestContext
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.policies.EmptyPolicy
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.policies.RelayLimits
|
||||
import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
|
||||
@@ -44,7 +45,10 @@ class RelayLimitsServerTest {
|
||||
|
||||
private val emptySource =
|
||||
object : EventSource {
|
||||
override fun events(filters: List<Filter>): Flow<Event> = emptyFlow()
|
||||
override fun events(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
): Flow<Event> = emptyFlow()
|
||||
}
|
||||
|
||||
private class Collector {
|
||||
@@ -79,7 +83,10 @@ class RelayLimitsServerTest {
|
||||
// An empty flow EOSEs immediately; no need to keep it open.
|
||||
val source =
|
||||
object : EventSource {
|
||||
override fun events(filters: List<Filter>): Flow<Event> = emptyFlow()
|
||||
override fun events(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
): Flow<Event> = emptyFlow()
|
||||
}
|
||||
val server = EventSourceServer(source, parentContext = dispatcher, limits = RelayLimits(maxSubscriptions = 2))
|
||||
val collector = Collector()
|
||||
@@ -104,7 +111,10 @@ class RelayLimitsServerTest {
|
||||
val seen = mutableListOf<Int?>()
|
||||
val source =
|
||||
object : EventSource {
|
||||
override fun events(filters: List<Filter>): Flow<Event> {
|
||||
override fun events(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
): Flow<Event> {
|
||||
seen.add(filters.single().limit)
|
||||
return emptyFlow()
|
||||
}
|
||||
|
||||
+5
-1
@@ -23,6 +23,7 @@ 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.backend.EventSource
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.backend.RequestContext
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.policies.EmptyPolicy
|
||||
import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
|
||||
import kotlinx.coroutines.ExperimentalCoroutinesApi
|
||||
@@ -51,7 +52,10 @@ class RelayServerListenerTest {
|
||||
|
||||
private val emptySource =
|
||||
object : EventSource {
|
||||
override fun events(filters: List<Filter>): Flow<Event> = emptyFlow()
|
||||
override fun events(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
): Flow<Event> = emptyFlow()
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
Reference in New Issue
Block a user