feat(nip67): EOSE completeness hints (finish / more / auth)

Wire model:
- EoseMessage gains optional `hints` (the third EOSE element) with
  isFinished()/hasMore()/needsAuth(); both the kotlinx and Jackson parsers
  read it (non-string/unknown entries ignored, non-array third element
  ignored) and both serializers write it only when present. The two-element
  fast path of toJson() is unchanged.

Client:
- SubscriptionListener gets onEose(relay, forFilters, hints); the pool calls
  it and the default forwards to the old two-argument onEose.
- fetchAllPages: "finish" ends the walk without the extra empty-page REQ
  (DRAINED, or LIMIT_REACHED/UNPAGEABLE when a filter already dropped out);
  "more" keeps paging; "auth" waits once for the NIP-42 verdict like an
  auth-required CLOSED and reads the post-AUTH re-served page, dropping ids
  it already delivered. An unanswered "auth" can never yield DRAINED.
- RelayAuthenticator treats an EOSE "auth" hint like an auth-required
  CLOSED (non-interactive re-auth on the stored challenge, whose OK re-REQs
  via syncFilters), at most once per challenge.

Relay engine (opt-in, RelayServerBase.completenessHints, default off):
- One filter with limit L: query L+1, forward L, send "more" if the extra
  row exists else "finish". All filters unbounded: "finish". Several
  filters: "finish" only when fewer rows than the smallest limit; else no
  hint. limit 0 and the policy-screened path never get a hint. Opt-in
  because it relies on the backend honouring limit exactly.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01MGR1u8SyzcUuekub39SBsc
This commit is contained in:
Claude
2026-09-27 19:28:38 +00:00
parent 84bc404dfe
commit 00eb7cf447
14 changed files with 855 additions and 15 deletions
@@ -37,6 +37,7 @@ import kotlinx.serialization.descriptors.SerialDescriptor
import kotlinx.serialization.descriptors.buildClassSerialDescriptor
import kotlinx.serialization.encoding.Decoder
import kotlinx.serialization.encoding.Encoder
import kotlinx.serialization.json.JsonArray
import kotlinx.serialization.json.JsonDecoder
import kotlinx.serialization.json.JsonEncoder
import kotlinx.serialization.json.JsonObject
@@ -98,6 +99,10 @@ object MessageKSerializer : KSerializer<Message> {
is EoseMessage -> {
add(JsonPrimitive(value.subId))
// NIP-67: optional third element, the completeness hints.
value.hints?.let { hints ->
add(buildJsonArray { hints.forEach { add(JsonPrimitive(it)) } })
}
}
is LimitsMessage -> {
@@ -136,7 +141,12 @@ object MessageKSerializer : KSerializer<Message> {
}
EoseMessage.LABEL -> {
EoseMessage(array[1].jsonPrimitive.content)
// NIP-67: an optional array of hint strings; anything else there is ignored.
val hints =
(array.getOrNull(2) as? JsonArray)?.mapNotNull { hint ->
(hint as? JsonPrimitive)?.takeIf { it.isString }?.content
}
EoseMessage(array[1].jsonPrimitive.content, hints)
}
NoticeMessage.LABEL -> {
@@ -30,6 +30,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.client.auth.awaitAuthOutcome
import com.vitorpamplona.quartz.nip01Core.relay.client.auth.hasAuthResponder
import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener
import com.vitorpamplona.quartz.nip01Core.relay.client.single.newSubId
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.EoseMessage
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.MachineReadablePrefix
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
@@ -190,6 +191,23 @@ data class PagedFetchResult(
* it a `limit` to bound that single page; without one you get the relay's default page
* of top hits.
*
* **NIP-67 completeness hints.** A relay may append hints to a page's `EOSE`:
*
* - `"finish"` — every stored match was sent, so the walk stops right there instead of
* spending one more REQ just to observe an empty page. It ends
* [PagedFetchResult.End.DRAINED] (or LIMIT_REACHED / UNPAGEABLE when a filter had
* already dropped out of the page, since `finish` can only speak for what was asked).
* - `"more"` — the relay holds more; paging continues, which is what the walk does
* anyway until it sees an empty page, so this needs no special handling.
* - `"auth"` — more may be available after NIP-42. Handled like an `auth-required:`
* CLOSED: when a responder is attached the walk waits (once) for the AUTH verdict and,
* on success, reads the page the relay re-serves after the AUTH's re-REQ, dropping the
* events it already delivered. If the AUTH does not happen, the page's events still
* count but the walk can no longer claim DRAINED: it ends
* [PagedFetchResult.End.AUTH_REQUIRED] wherever it would have ended DRAINED.
*
* Hints are only ever a shortcut; their absence changes nothing (the heuristic above).
*
* @param relay The relay to query.
* @param filters Filters to apply on every page (the `until` field is overwritten per page).
* @param idleTimeoutMs Idle window per page — like every accessory timeout, it is measured
@@ -336,6 +354,16 @@ suspend fun INostrClient.fetchAllPages(
// declining to give one.
var pageEnd: PageSignal? = null
// NIP-67 hints of the EOSE that ended this page (null: none sent). Written on the
// relay's reader thread before the EOSE signal is sent; the channel orders it.
var eoseHints: List<String>? = null
// Ids delivered on this page, kept only while an EOSE `"auth"` hint could still make
// the relay re-serve the page after AUTH (at most once per walk), so the re-served
// copies of events already handed to [onEvent] are dropped. Reader-thread only.
val pageIds: HashSet<HexKey>? = if (pendingOnAuthRequired && !authRetried) HashSet() else null
var reServing = false
try {
val listener =
object : SubscriptionListener {
@@ -363,6 +391,9 @@ suspend fun INostrClient.fetchAllPages(
// Drop a boundary-second event we already delivered on an
// earlier page (the inclusive re-fetch returns it again).
if (boundary != null && event.createdAt == boundary && event.id in seenAtBoundary) return
// The relay re-serving this page after an EOSE "auth" hint: skip what
// this page already delivered.
if (reServing && pageIds != null && event.id in pageIds) return
// Count this event against every active filter it satisfies
// (one event can match more than one). Only a non-search filter
@@ -390,6 +421,7 @@ suspend fun INostrClient.fetchAllPages(
if (atLeastOne) {
onEvent(event)
delivered++
pageIds?.add(event.id)
// Track the oldest advancing second and the ids delivered
// in it — that becomes the next boundary and its dedup set.
if (advancesCursor) {
@@ -414,6 +446,15 @@ suspend fun INostrClient.fetchAllPages(
doneChannel.trySend(PageSignal.EOSE)
}
override fun onEose(
relay: NormalizedRelayUrl,
forFilters: List<Filter>?,
hints: List<String>?,
) {
eoseHints = hints
doneChannel.trySend(PageSignal.EOSE)
}
override fun onClosed(
message: String,
relay: NormalizedRelayUrl,
@@ -454,6 +495,18 @@ suspend fun INostrClient.fetchAllPages(
clock.bump()
pageEnd = doneChannel.receiveWithinIdle(clock, idleTimeoutMs)
}
} else if (pageEnd == PageSignal.EOSE && eoseHints.hasHint(EoseMessage.HINT_AUTH) && pendingOnAuthRequired && !authRetried) {
// NIP-67 "auth": the page was answered, but the relay says it held some back.
// Same wait as the CLOSED case — the relay sent its challenge before this EOSE,
// and the AUTH's OK re-sends this very REQ — except the page already delivered
// events, so the re-served copies are dropped via [pageIds].
authRetried = true
reServing = true
if (awaitAuthOutcome(relay, authMark, DEFAULT_AUTH_GRACE_MS, idleTimeoutMs) == AuthOutcome.AUTHENTICATED) {
eoseHints = null
clock.bump()
pageEnd = doneChannel.receiveWithinIdle(clock, idleTimeoutMs)
}
}
unsubscribe(subId)
@@ -465,6 +518,10 @@ suspend fun INostrClient.fetchAllPages(
totalEvents += delivered
// The page ended on an EOSE saying more is visible only after AUTH, and no AUTH
// took the wall down: whatever it did deliver stands, but it cannot prove absence.
val authBlocked = pageEnd == PageSignal.EOSE && eoseHints.hasHint(EoseMessage.HINT_AUTH)
// The relay sent nothing at-or-below `until`. Whether that DRAINS the set
// depends on why the page ended and on what was asked:
//
@@ -490,6 +547,23 @@ suspend fun INostrClient.fetchAllPages(
pageEnd == null -> PagedFetchResult.End.IDLE
cappedByLimit -> PagedFetchResult.End.LIMIT_REACHED
filters.any { it.search != null } -> PagedFetchResult.End.UNPAGEABLE
authBlocked -> PagedFetchResult.End.AUTH_REQUIRED
else -> PagedFetchResult.End.DRAINED
}
break
}
// NIP-67 "finish": the relay says it sent every stored match for this page's
// filters, so there is nothing below the cursor to ask for — stop now rather than
// spend a REQ to watch an empty page come back. It only speaks for the filters this
// page actually carried; one that already dropped out (limit met, or a search after
// its single page) keeps the reading it would have had.
if (pageEnd == PageSignal.EOSE && eoseHints.hasHint(EoseMessage.HINT_FINISH)) {
end =
when {
authBlocked -> PagedFetchResult.End.AUTH_REQUIRED
filters.indices.any { i -> filters[i].limit.let { it != null && matchCountPerFilter[i] >= it } } -> PagedFetchResult.End.LIMIT_REACHED
filters.any { it.search != null } -> PagedFetchResult.End.UNPAGEABLE
else -> PagedFetchResult.End.DRAINED
}
break
@@ -587,3 +661,5 @@ suspend fun INostrClient.fetchAllPages(
onNewPage = onNewPage,
onEvent = onEvent,
)
private fun List<String>?.hasHint(hint: String) = this != null && contains(hint)
@@ -25,6 +25,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.client.listeners.RelayConnection
import com.vitorpamplona.quartz.nip01Core.relay.client.single.IRelayClient
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.AuthMessage
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.ClosedMessage
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.EoseMessage
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.MachineReadablePrefix
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.OkMessage
@@ -113,6 +114,9 @@ class RelayAuthenticator(
// from RelayAuthStatus.snapshot().
private val authStatus = LargeCache<NormalizedRelayUrl, RelayAuthStatus>()
/** The challenge each relay was last re-authenticated on because of an EOSE `"auth"` hint. */
private val authHintRetried = LargeCache<NormalizedRelayUrl, String>()
private val _authStateFlow = MutableStateFlow<PersistentMap<NormalizedRelayUrl, RelayAuthSnapshot>>(persistentMapOf())
/**
@@ -143,6 +147,7 @@ class RelayAuthenticator(
is AuthMessage -> authenticate(relay, msg.challenge, interactive = true)
is OkMessage -> checkAuthResults(relay, msg)
is ClosedMessage -> reauthenticateIfAuthRequired(relay, msg)
is EoseMessage -> reauthenticateIfAuthHinted(relay, msg)
}
}
@@ -153,6 +158,7 @@ class RelayAuthenticator(
override fun onDisconnected(relay: IRelayClient) {
authStatus.remove(relay.url)
authHintRetried.remove(relay.url)
publishSnapshot(relay.url)
}
}
@@ -222,6 +228,34 @@ class RelayAuthenticator(
msg: ClosedMessage,
) {
if (MachineReadablePrefix.parse(msg.message) != MachineReadablePrefix.AUTH_REQUIRED) return
reauthenticateWithStoredChallenge(relay)
}
/**
* NIP-67 / NIP-42: an `EOSE` carrying the `"auth"` hint says the relay may hold more
* matches for this subscription if we authenticate. The relay MUST have sent its
* `AUTH` challenge before that EOSE, so the challenge is already stored and the normal
* [authenticate] pass has usually run on it. This takes the same path as an
* `auth-required:` CLOSED: re-attach any approved identity not yet sent on that
* challenge (never prompting), and let the AUTH's `OK` → [INostrClient.syncFilters]
* re-send the REQ so the relay can serve what it held back. Deduped per
* (pubkey, challenge) and skipped while an AUTH is in flight, so it cannot loop.
*/
private fun reauthenticateIfAuthHinted(
relay: IRelayClient,
msg: EoseMessage,
) {
if (!msg.needsAuth()) return
// A relay that keeps refusing us may tag EVERY EOSE with "auth". Unlike a CLOSED, the
// subscription is still answered, so there is no refusal to recover from — one retry
// per challenge is enough, and it spares an external signer a pass per subscription.
val challenge = authStatus.get(relay.url)?.lastChallenge() ?: return
if (authHintRetried.get(relay.url) == challenge) return
authHintRetried.put(relay.url, challenge)
reauthenticateWithStoredChallenge(relay)
}
private fun reauthenticateWithStoredChallenge(relay: IRelayClient) {
val status = authStatus.get(relay.url) ?: return
// Coalesce the burst: a relay refuses EVERY currently-open sub with its own `auth-required`
// CLOSED, so a single missing identity yields many CLOSEDs at once. Re-signing on each would
@@ -302,6 +302,7 @@ class PoolRequests(
desiredSubListeners.get(msg.subId)?.onEose(
relay = relay.url,
forFilters = forFilters,
hints = msg.hints,
)
// send a newer version when done
@@ -30,6 +30,18 @@ interface SubscriptionListener {
forFilters: List<Filter>?,
) {}
/**
* EOSE together with its NIP-67 completeness [hints] (`finish`, `more`, `auth`, …;
* null when the relay sent the plain two-element EOSE). The pool calls this one;
* the default forwards to the two-argument [onEose], so listeners that don't care
* about hints keep overriding that.
*/
fun onEose(
relay: NormalizedRelayUrl,
forFilters: List<Filter>?,
hints: List<String>?,
) = onEose(relay, forFilters)
suspend fun onEvent(
event: Event,
isLive: Boolean,
@@ -20,11 +20,29 @@
*/
package com.vitorpamplona.quartz.nip01Core.relay.commands.toClient
/**
* `["EOSE", <subId>]`, optionally with NIP-67 completeness hints:
* `["EOSE", <subId>, [<hint>, ...]]`.
*
* [hints] is null when the relay sent the two-element form. Hints are only
* about stored events; their presence is definitive, their absence is not
* (see [isFinished] / [hasMore]). Unknown hint values are kept but ignored.
*/
class EoseMessage(
val subId: String,
val hints: List<String>? = null,
) : Message {
override fun label() = LABEL
/** NIP-67 `finish`: every stored match was sent; do not paginate further. */
fun isFinished() = hints?.contains(HINT_FINISH) == true
/** NIP-67 `more`: the relay holds more stored matches than it sent; paginate. */
fun hasMore() = hints?.contains(HINT_MORE) == true
/** NIP-67 `auth`: more stored matches may be available after NIP-42 AUTH. */
fun needsAuth() = hints?.contains(HINT_AUTH) == true
/**
* Wire form is `["EOSE","<subId>"]` — sent once per REQ, so it is on
* the per-subscription floor. Splice it directly when [subId] needs no
@@ -33,7 +51,7 @@ class EoseMessage(
* any exotic subId falls back.
*/
override fun toJson(): String {
if (!isEscapeFreeAscii(subId)) return super.toJson()
if (hints != null || !isEscapeFreeAscii(subId)) return super.toJson()
return buildString(subId.length + 12) {
append("[\"EOSE\",\"")
append(subId)
@@ -43,5 +61,10 @@ class EoseMessage(
companion object {
const val LABEL = "EOSE"
// NIP-67 hint values.
const val HINT_FINISH = "finish"
const val HINT_MORE = "more"
const val HINT_AUTH = "auth"
}
}
@@ -0,0 +1,99 @@
/*
* 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.commands.toClient.EoseMessage
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
/**
* Works out the NIP-67 completeness hint for one REQ's stored replay, only ever
* claiming what the store's answer actually proves:
*
* - **One filter with `limit = L > 0`**: the store is asked for `L + 1` rows and the
* extra (oldest) row is withheld. Seeing it proves the relay holds more (`"more"`);
* not seeing it proves the replay was complete (`"finish"`). The client receives
* exactly the `L` events it would have without the probe.
* - **Every filter unbounded (`limit = null`)**: the store returns every match
* (STORE-F12), so the replay is complete (`"finish"`).
* - **Several filters, some limited**: limits are per filter and the replay is their
* deduped union, so rows cannot be attributed back to a filter without matching
* each one. Only the cheap, sound case is claimed: fewer rows than the smallest
* limit means no filter reached its limit (`"finish"`). Anything else sends no hint,
* which NIP-67 allows (absence is not definitive).
*
* `limit = 0` never gets a hint and is never probed: NIP-01 forbids returning stored
* events for it, so it keeps its exact query.
*
* This assumes the backing store honours `limit` exactly and treats `null` as
* unbounded, which is why [RelaySession] only uses it when the server opts in.
*/
internal class EoseCompletenessProbe private constructor(
/** The filters to actually query (the limit may be raised by one). */
val queryFilters: List<Filter>,
/** Max stored events to forward; the rest are only counted. Null: forward all. */
private val forwardCap: Int?,
private val rule: Rule,
) {
private enum class Rule { PROBE, ALL_UNBOUNDED, UNDER_MIN_LIMIT }
private val minLimit: Int = if (rule == Rule.UNDER_MIN_LIMIT) queryFilters.minOf { it.limit ?: Int.MAX_VALUE } else 0
/** Stored events the store produced for this REQ, forwarded or not. */
private var seen = 0
/**
* Counts one stored event and says whether it should be forwarded to the client.
* Called from the single replay coroutine, before EOSE.
*/
fun onStored(): Boolean {
seen++
return forwardCap == null || seen <= forwardCap
}
/** The hints for this replay's EOSE, or null for none. */
fun hints(): List<String>? =
when (rule) {
Rule.PROBE -> if (seen > forwardCap!!) MORE else FINISH
Rule.ALL_UNBOUNDED -> FINISH
Rule.UNDER_MIN_LIMIT -> if (seen < minLimit) FINISH else null
}
companion object {
private val FINISH = listOf(EoseMessage.HINT_FINISH)
private val MORE = listOf(EoseMessage.HINT_MORE)
fun of(filters: List<Filter>): EoseCompletenessProbe? {
if (filters.isEmpty()) return null
if (filters.any { it.limit == 0 }) return null
if (filters.size == 1) {
val limit = filters[0].limit
if (limit != null && limit < Int.MAX_VALUE) {
return EoseCompletenessProbe(listOf(filters[0].copy(limit = limit + 1)), limit, Rule.PROBE)
}
}
if (filters.all { it.limit == null }) return EoseCompletenessProbe(filters, null, Rule.ALL_UNBOUNDED)
return EoseCompletenessProbe(filters, null, Rule.UNDER_MIN_LIMIT)
}
}
}
@@ -63,6 +63,15 @@ abstract class RelayServerBase(
/** Number of connections currently registered with this server. */
val activeConnections: Long get() = connections.active
/**
* NIP-67 opt-in: when true, connections opened from now on append `"finish"` /
* `"more"` to their EOSEs where the stored replay proves it (see
* [RelaySession.completenessHints]). Enable it only for a backend that honours
* `limit` exactly and returns every match for an unbounded filter, and advertise
* `67` in the NIP-11 `supported_nips` when you do.
*/
var completenessHints: Boolean = false
/**
* Builds the per-connection policy, prepending a [LimitsPolicy] when
* [limits] is set so requests are clamped/rejected before the application
@@ -91,6 +100,7 @@ abstract class RelayServerBase(
sink = sink,
onClose = { connections.unregister(it.id) },
negentropySettings = negentropySettings,
completenessHints = completenessHints,
),
)
@@ -79,6 +79,13 @@ class RelaySession(
* open/close of the same connection. Defaults to a fresh monotonic id.
*/
val id: Long = nextConnectionId(),
/**
* NIP-67: append a completeness hint (`"finish"` / `"more"`) to each REQ's `EOSE`
* when the stored replay proves one — see [EoseCompletenessProbe]. Off by default:
* the proof assumes the [store] honours `limit` exactly and returns every match
* for an unbounded filter, which an arbitrary backend need not do.
*/
val completenessHints: Boolean = false,
) : AutoCloseable {
/** The original, string-only constructor; every frame goes to [onSend] as wire JSON. */
constructor(
@@ -353,7 +360,15 @@ class RelaySession(
}
// Policy may rewrite filters to match the user's access level.
val filters = (result as PolicyResult.Accepted).cmd.filters
val acceptedFilters = (result as PolicyResult.Accepted).cmd.filters
// NIP-67: may raise a single filter's limit by one to detect "more"; the extra
// stored row is counted but never sent. Zero-decode path only: the screened path's
// single `onEach` also carries live events accepted mid-replay, so stored rows can't
// be told apart there, and a policy that vetoes rows could not honestly say "finish".
val probe = if (completenessHints && !policy.filtersOutgoingEvents) EoseCompletenessProbe.of(acceptedFilters) else null
val filters = probe?.queryFilters ?: acceptedFilters
val eose = { send(EoseMessage(cmd.subId, probe?.hints())) }
// UNDISPATCHED: the stored replay runs inline on this coroutine —
// the reader-pool acquire doesn't suspend when a connection is
@@ -377,7 +392,7 @@ class RelaySession(
send(EventMessage(cmd.subId, event))
}
},
onEose = { send(EoseMessage(cmd.subId)) },
onEose = { eose() },
)
} else {
// Zero-decode path: the stored replay splices raw
@@ -395,13 +410,16 @@ class RelaySession(
ctx = requestContext,
filters = filters,
onEachStored = { raw ->
sendRaw(
buildString(framePrefix.length + raw.jsonTags.length + raw.content.length + 256) {
append(framePrefix)
raw.appendJsonObjectTo(this)
append(']')
},
)
// NIP-67 probe: the one extra row it asked for is counted, not sent.
if (probe == null || probe.onStored()) {
sendRaw(
buildString(framePrefix.length + raw.jsonTags.length + raw.content.length + 256) {
append(framePrefix)
raw.appendJsonObjectTo(this)
append(']')
},
)
}
},
// Live events arrive with their wire body already
// serialized (once per event, shared across every
@@ -417,7 +435,7 @@ class RelaySession(
},
)
},
onEose = { send(EoseMessage(cmd.subId)) },
onEose = { eose() },
)
}
} catch (e: CancellationException) {
@@ -0,0 +1,198 @@
/*
* Copyright (c) 2025 Vitor Pamplona
*
* Permission is hereby granted, free of charge, to any person obtaining a copy of
* this software and associated documentation files (the "Software"), to deal in
* the Software without restriction, including without limitation the rights to use,
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
* Software, and to permit persons to whom the Software is furnished to do so,
* subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all
* copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
*/
package com.vitorpamplona.quartz.nip01Core.relay.server
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper
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.toClient.Message
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.EventCmd
import com.vitorpamplona.quartz.nip01Core.relay.server.policies.EmptyPolicy
import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
import kotlinx.coroutines.CoroutineDispatcher
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.test.UnconfinedTestDispatcher
import kotlinx.coroutines.test.runTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNull
import kotlin.test.assertTrue
/** The relay engine's opt-in NIP-67 EOSE completeness hints. */
@OptIn(ExperimentalCoroutinesApi::class)
class NostrServerEoseTest {
private val pubkey = "46fcbe3065eaf1ae7811465924e48923363ff3f526bd6f73d7c184b16bd8ce4d"
private val sig = "4aa5264965018fa12a326686ad3d3bd8beae3218dcc83689b19ca1e6baeb791531943c15363aa6707c7c0c8b2d601deca1f20c32078b2872d356cdca03b04cce"
private val noop: (String) -> Unit = {}
private fun hexId(n: Int): String = n.toString().padStart(64, '0')
private fun testEvent(
id: String,
kind: Int = 1,
createdAt: Long = 1000L,
) = Event(id, pubkey, createdAt, kind, emptyArray(), "hello", sig)
private fun server(
dispatcher: CoroutineDispatcher,
store: EventStore,
hints: Boolean = false,
) = NostrServer(store = store, policyBuilder = { EmptyPolicy }, parentContext = dispatcher).also {
it.completenessHints = hints
}
private class Collector {
val messages = mutableListOf<String>()
val send: (String) -> Unit = { messages.add(it) }
fun parsed(): List<Message> =
messages
.filter { it.startsWith("[\"EVENT\"") || it.startsWith("[\"EOSE\"") }
.map { OptimizedJsonMapper.fromJsonToMessage(it) }
fun events() = parsed().filterIsInstance<EventMessage>()
fun eoses() = parsed().filterIsInstance<EoseMessage>()
}
private suspend fun EventStore.seed(count: Int) {
for (i in 1..count) insert(testEvent(hexId(i), createdAt = i.toLong()))
}
// -- NIP-67: completeness hints ------------------------------------------------
@Test
fun hintsAreOffByDefault() =
runTest {
val dispatcher = UnconfinedTestDispatcher(testScheduler)
val store = EventStore(null)
store.seed(3)
val server = server(dispatcher, store)
val collector = Collector()
server.connect(collector.send).receive("""["REQ","s",{"kinds":[1]}]""")
assertTrue(collector.messages.contains("""["EOSE","s"]"""))
assertNull(collector.eoses().single().hints)
server.close()
}
@Test
fun truncatedByLimitSendsMore() =
runTest {
val dispatcher = UnconfinedTestDispatcher(testScheduler)
val store = EventStore(null)
store.seed(10)
val server = server(dispatcher, store, hints = true)
val collector = Collector()
server.connect(collector.send).receive("""["REQ","s",{"kinds":[1],"limit":3}]""")
// The probe asks the store for 4 but only the newest 3 go out.
assertEquals(listOf(hexId(10), hexId(9), hexId(8)), collector.events().map { it.event.id })
assertEquals(listOf(EoseMessage.HINT_MORE), collector.eoses().single().hints)
server.close()
}
@Test
fun exactlyTheLimitSendsFinish() =
runTest {
val dispatcher = UnconfinedTestDispatcher(testScheduler)
val store = EventStore(null)
store.seed(10)
val server = server(dispatcher, store, hints = true)
val collector = Collector()
server.connect(collector.send).receive("""["REQ","s",{"kinds":[1],"limit":10}]""")
assertEquals(10, collector.events().size)
assertEquals(listOf(EoseMessage.HINT_FINISH), collector.eoses().single().hints)
server.close()
}
@Test
fun unboundedFilterSendsFinish() =
runTest {
val dispatcher = UnconfinedTestDispatcher(testScheduler)
val store = EventStore(null)
store.seed(4)
val server = server(dispatcher, store, hints = true)
val collector = Collector()
server.connect(collector.send).receive("""["REQ","s",{"kinds":[1]}]""")
assertEquals(4, collector.events().size)
assertEquals(listOf(EoseMessage.HINT_FINISH), collector.eoses().single().hints)
server.close()
}
@Test
fun multiFilterHintsOnlyWhatTheCountProves() =
runTest {
val dispatcher = UnconfinedTestDispatcher(testScheduler)
val store = EventStore(null)
store.seed(4)
val server = server(dispatcher, store, hints = true)
val under = Collector()
server.connect(under.send).receive("""["REQ","s",{"kinds":[1],"limit":10},{"kinds":[2],"limit":20}]""")
assertEquals(4, under.events().size)
assertEquals(listOf(EoseMessage.HINT_FINISH), under.eoses().single().hints, "4 rows < smallest limit: nothing was cut")
val ambiguous = Collector()
server.connect(ambiguous.send).receive("""["REQ","s",{"kinds":[1],"limit":2},{"kinds":[2],"limit":20}]""")
assertEquals(2, ambiguous.events().size)
assertNull(ambiguous.eoses().single().hints, "cannot attribute rows to filters cheaply, so no claim")
server.close()
}
@Test
fun limitZeroNeverGetsAHint() =
runTest {
val dispatcher = UnconfinedTestDispatcher(testScheduler)
val store = EventStore(null)
store.seed(4)
val server = server(dispatcher, store, hints = true)
val collector = Collector()
server.connect(collector.send).receive("""["REQ","s",{"kinds":[1],"limit":0}]""")
assertEquals(0, collector.events().size)
assertNull(collector.eoses().single().hints)
server.close()
}
@Test
fun probeDoesNotHoldBackLiveEvents() =
runTest {
val dispatcher = UnconfinedTestDispatcher(testScheduler)
val store = EventStore(null)
store.seed(5)
val server = server(dispatcher, store, hints = true)
val collector = Collector()
server.connect(collector.send).receive("""["REQ","s",{"kinds":[1],"limit":1}]""")
assertEquals(1, collector.events().size)
server.connect(noop).receive(OptimizedJsonMapper.toJson(EventCmd(testEvent(hexId(100), createdAt = 9000L))))
server.connect(noop).receive(OptimizedJsonMapper.toJson(EventCmd(testEvent(hexId(101), createdAt = 9001L))))
assertEquals(listOf(hexId(5), hexId(100), hexId(101)), collector.events().map { it.event.id })
server.close()
}
}
@@ -54,9 +54,26 @@ class MessageDeserializer : StdDeserializer<Message>(Message::class.java) {
}
EoseMessage.LABEL -> {
EoseMessage(
subId = jp.nextTextValue(),
)
val subId = jp.nextTextValue()
// NIP-67: an optional third element, an array of hint strings. The array is
// consumed here (stepping past its END_ARRAY) so the drain loop below only
// ever sees the outer frame's tokens.
val hints =
if (jp.nextToken() == JsonToken.START_ARRAY) {
val list = ArrayList<String>(2)
while (jp.nextToken() != JsonToken.END_ARRAY) {
if (jp.currentToken == JsonToken.VALUE_STRING) {
list.add(jp.text)
} else {
jp.skipChildren()
}
}
jp.nextToken()
list
} else {
null
}
EoseMessage(subId, hints)
}
NoticeMessage.LABEL -> {
@@ -78,6 +78,12 @@ class MessageSerializer : StdSerializer<Message>(Message::class.java) {
is EoseMessage -> {
gen.writeString(msg.subId)
// NIP-67: optional third element, the completeness hints.
msg.hints?.let { hints ->
gen.writeStartArray()
hints.forEach { gen.writeString(it) }
gen.writeEndArray()
}
}
is LimitsMessage -> {
@@ -0,0 +1,201 @@
/*
* 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
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.relay.client.EmptyNostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.PagedFetchResult
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAllPages
import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import kotlin.test.Test
import kotlin.test.assertEquals
/**
* NIP-67 completeness hints steering [fetchAllPages]: `finish` ends the walk without the
* extra empty-page REQ, `more` keeps paging, and an unanswered `auth` hint keeps the walk
* from claiming DRAINED.
*/
class NostrClientFetchAllPagesEoseHintsTest {
private class ScriptedClient : INostrClient by EmptyNostrClient() {
@Volatile
var listener: SubscriptionListener? = null
@Volatile
var subscribeCount = 0
override fun subscribe(
subId: String,
filters: Map<NormalizedRelayUrl, List<Filter>>,
listener: SubscriptionListener?,
) {
subscribeCount++
this.listener = listener
}
suspend fun awaitPage(n: Int) {
while (subscribeCount < n) delay(2)
}
}
private val relay = RelayUrlNormalizer.normalize("wss://hints.example.com")
private fun event(createdAt: Long) =
Event(
id = createdAt.toString(16).padStart(64, '0'),
pubKey = "f".repeat(64),
createdAt = createdAt,
kind = 1,
tags = emptyArray(),
content = "e$createdAt",
sig = "0".repeat(128),
)
@Test
fun finishEndsTheWalkWithoutAnotherPage() =
runBlocking {
val client = ScriptedClient()
val feeder =
launch {
client.awaitPage(1)
client.listener!!.onEvent(event(2000), false, relay, null)
client.listener!!.onEvent(event(1000), false, relay, null)
client.listener!!.onEose(relay, null, listOf("finish"))
}
val result = client.fetchAllPages(relay = relay, filters = listOf(Filter(kinds = listOf(1))), idleTimeoutMs = 2_000) { }
feeder.join()
assertEquals(2, result.downloaded)
assertEquals(PagedFetchResult.End.DRAINED, result.end, "the relay said it sent every stored match")
assertEquals(1, client.subscribeCount, "no second REQ just to observe an empty page")
}
@Test
fun moreKeepsPaging() =
runBlocking {
val client = ScriptedClient()
val feeder =
launch {
client.awaitPage(1)
client.listener!!.onEvent(event(2000), false, relay, null)
client.listener!!.onEose(relay, null, listOf("more"))
client.awaitPage(2)
client.listener!!.onEvent(event(1000), false, relay, null)
client.listener!!.onEose(relay, null, listOf("finish"))
}
val result = client.fetchAllPages(relay = relay, filters = listOf(Filter(kinds = listOf(1))), idleTimeoutMs = 2_000) { }
feeder.join()
assertEquals(2, result.downloaded)
assertEquals(PagedFetchResult.End.DRAINED, result.end)
assertEquals(2, client.subscribeCount)
}
@Test
fun unknownHintsAreIgnored() =
runBlocking {
val client = ScriptedClient()
val feeder =
launch {
client.awaitPage(1)
client.listener!!.onEvent(event(2000), false, relay, null)
client.listener!!.onEose(relay, null, listOf("somethingNew"))
client.awaitPage(2)
client.listener!!.onEose(relay, null, emptyList())
}
val result = client.fetchAllPages(relay = relay, filters = listOf(Filter(kinds = listOf(1))), idleTimeoutMs = 2_000) { }
feeder.join()
assertEquals(1, result.downloaded)
assertEquals(PagedFetchResult.End.DRAINED, result.end, "falls back to the empty-page heuristic")
assertEquals(2, client.subscribeCount)
}
@Test
fun finishWithAnUnansweredAuthHintCannotClaimDrained() =
runBlocking {
// No NIP-42 responder is attached, so the "auth" hint cannot be acted on.
val client = ScriptedClient()
val feeder =
launch {
client.awaitPage(1)
client.listener!!.onEvent(event(2000), false, relay, null)
client.listener!!.onEose(relay, null, listOf("auth", "finish"))
}
val result = client.fetchAllPages(relay = relay, filters = listOf(Filter(kinds = listOf(1))), idleTimeoutMs = 2_000) { }
feeder.join()
assertEquals(1, result.downloaded, "what was delivered still counts")
assertEquals(PagedFetchResult.End.AUTH_REQUIRED, result.end, "the relay may hold more for an authenticated user")
assertEquals(1, client.subscribeCount)
}
@Test
fun anEmptyPageWithAnAuthHintIsNotADrain() =
runBlocking {
val client = ScriptedClient()
val feeder =
launch {
client.awaitPage(1)
client.listener!!.onEvent(event(2000), false, relay, null)
client.listener!!.onEose(relay, null, listOf("auth"))
client.awaitPage(2)
client.listener!!.onEose(relay, null, listOf("auth"))
}
val result = client.fetchAllPages(relay = relay, filters = listOf(Filter(kinds = listOf(1))), idleTimeoutMs = 2_000) { }
feeder.join()
assertEquals(1, result.downloaded)
assertEquals(PagedFetchResult.End.AUTH_REQUIRED, result.end)
}
@Test
fun finishAfterAFilterMetItsLimitReportsLimitReached() =
runBlocking {
val client = ScriptedClient()
val feeder =
launch {
client.awaitPage(1)
client.listener!!.onEvent(event(2000), false, relay, null)
client.listener!!.onEvent(event(1000), false, relay, null)
client.listener!!.onEose(relay, null, listOf("finish"))
}
val result = client.fetchAllPages(relay = relay, filters = listOf(Filter(kinds = listOf(1), limit = 2)), idleTimeoutMs = 2_000) { }
feeder.join()
assertEquals(2, result.downloaded)
assertEquals(PagedFetchResult.End.LIMIT_REACHED, result.end)
assertEquals(1, client.subscribeCount)
}
}
@@ -0,0 +1,135 @@
/*
* 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.commands.toClient
import com.vitorpamplona.quartz.nip01Core.jackson.JacksonMapper
import com.vitorpamplona.quartz.nip01Core.kotlinSerialization.KotlinSerializationMapper
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertIs
import kotlin.test.assertNull
import kotlin.test.assertTrue
/** NIP-67 EOSE completeness hints, through both JSON backends. */
class EoseHintsParsingTest {
private val parsers: List<Pair<String, (String) -> Message>> =
listOf(
"jackson" to { json -> JacksonMapper.fromJsonToMessage(json) },
"kotlinx" to { json -> KotlinSerializationMapper.fromJsonToMessage(json) },
)
private fun eachParser(
json: String,
check: (String, EoseMessage) -> Unit,
) = parsers.forEach { (name, parse) ->
val msg = parse(json)
assertIs<EoseMessage>(msg, name)
check(name, msg)
}
@Test
fun twoElementEoseHasNoHints() =
eachParser("""["EOSE","sub1"]""") { name, msg ->
assertEquals("sub1", msg.subId, name)
assertNull(msg.hints, name)
assertFalse(msg.isFinished(), name)
assertFalse(msg.hasMore(), name)
assertFalse(msg.needsAuth(), name)
}
@Test
fun finishHint() =
eachParser("""["EOSE","sub2",["finish"]]""") { name, msg ->
assertEquals("sub2", msg.subId, name)
assertEquals(listOf("finish"), msg.hints, name)
assertTrue(msg.isFinished(), name)
assertFalse(msg.hasMore(), name)
}
@Test
fun moreHint() =
eachParser("""["EOSE","sub2b",["more"]]""") { name, msg ->
assertTrue(msg.hasMore(), name)
assertFalse(msg.isFinished(), name)
}
@Test
fun multipleAndUnknownHints() =
eachParser("""["EOSE","sub4",["auth","finish","somethingNew"]]""") { name, msg ->
assertEquals(listOf("auth", "finish", "somethingNew"), msg.hints, name)
assertTrue(msg.needsAuth(), name)
assertTrue(msg.isFinished(), name)
}
@Test
fun emptyHintArray() =
eachParser("""["EOSE","sub5",[]]""") { name, msg ->
assertEquals(emptyList(), msg.hints, name)
assertFalse(msg.isFinished(), name)
}
@Test
fun nonStringHintsAreIgnored() =
eachParser("""["EOSE","sub6",[1,{"a":[2]},["x"],"finish",null]]""") { name, msg ->
assertEquals(listOf("finish"), msg.hints, name)
}
@Test
fun nonArrayThirdElementIsIgnored() =
eachParser("""["EOSE","sub7","finish",{"x":1}]""") { name, msg ->
assertEquals("sub7", msg.subId, name)
assertNull(msg.hints, name)
}
@Test
fun trailingElementsAfterHintsAreTolerated() =
eachParser("""["EOSE","sub8",["more"],"extra",5]""") { name, msg ->
assertEquals(listOf("more"), msg.hints, name)
}
@Test
fun serializesWithoutHintsAsTwoElements() {
val msg = EoseMessage("sub1")
assertEquals("""["EOSE","sub1"]""", msg.toJson())
assertEquals("""["EOSE","sub1"]""", JacksonMapper.toJson(msg))
assertEquals("""["EOSE","sub1"]""", KotlinSerializationMapper.toJson(msg))
}
@Test
fun serializesHintsAsThirdElement() {
val msg = EoseMessage("sub1", listOf("auth", "finish"))
val expected = """["EOSE","sub1",["auth","finish"]]"""
assertEquals(expected, msg.toJson())
assertEquals(expected, JacksonMapper.toJson(msg))
assertEquals(expected, KotlinSerializationMapper.toJson(msg))
}
@Test
fun roundTripsThroughBothBackends() {
val json = EoseMessage("s", listOf("more")).toJson()
parsers.forEach { (name, parse) ->
val back = parse(json)
assertIs<EoseMessage>(back, name)
assertEquals(listOf("more"), back.hints, name)
}
}
}