From 4d5fe4bbdb83b70fdbe23c96031cd943bc7f9ffa Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 02:49:25 +0000 Subject: [PATCH 1/8] fix(quartz): NIP-98 replay cache refuses when full instead of forgetting a live token The verifier's replay cache evicted its oldest entry once it held 1024, live or not. On a public endpoint that is a replay primitive: sign 1024 throwaway tokens and a captured one inside its 60s window verifies again (reproduced against NIP-FE's /req). Entries now leave only by expiry; a cache full of live tokens answers the new one `rate-limited:` instead. The cap is a constructor parameter (default unchanged, 1024, sized for the admin rpc), so a public endpoint sizes it for its own signed-request rate. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_016jDNhr6a3J4VC5TaG3Yd48 --- .../quartz/nip98HttpAuth/Nip98AuthVerifier.kt | 31 ++++++++++--------- 1 file changed, 17 insertions(+), 14 deletions(-) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifier.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifier.kt index ebd580ee15..da83ca0d6d 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifier.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifier.kt @@ -52,10 +52,16 @@ class Nip98AuthVerifier( private val now: () -> Long = { TimeUtils.now() }, /** Allowed clock skew in seconds. NIP-98 says 60. */ private val toleranceSeconds: Long = 60, + /** + * Tokens remembered at once. When every remembered token is still inside its window, a new one + * is refused (`rate-limited:`) rather than an unexpired one forgotten: forgetting is what lets a + * captured token be replayed. Size it for the endpoint's signed-request rate over 2 x tolerance. + */ + private val maxReplayEntries: Int = MAX_REPLAY_ENTRIES, ) { /** * Recently-accepted event ids → expiry epoch second. Bounded to - * [MAX_REPLAY_ENTRIES] (insertion-order eviction); each entry expires + * [maxReplayEntries], never by evicting a live entry; each entry expires * after `2 × toleranceSeconds` (twice the accepted window so a token * can't be reused by an attacker who buffers across the boundary). * @@ -140,18 +146,13 @@ class Nip98AuthVerifier( while (it.hasNext()) { if (it.next().value <= nowSec) it.remove() else break } - if (seenEventIds.put(event.id, expiry) != null) { + if (seenEventIds.containsKey(event.id)) { return Result.Malformed("replay: this NIP-98 token has already been used") } - // Cap entries: drop oldest by insertion order. Equivalent to - // the JDK LinkedHashMap.removeEldestEntry hook we used before, - // but works in KMP commonMain. - while (seenEventIds.size > MAX_REPLAY_ENTRIES) { - val eldest = seenEventIds.keys.iterator() - if (!eldest.hasNext()) break - eldest.next() - eldest.remove() + if (seenEventIds.size >= maxReplayEntries) { + return Result.Malformed(REPLAY_CACHE_FULL) } + seenEventIds[event.id] = expiry } return Result.Verified(event.pubKey) @@ -173,11 +174,13 @@ class Nip98AuthVerifier( const val SCHEME = "Nostr " /** - * Cap on the in-memory replay-cache size. With a 60s tolerance - * an attacker would need to push >MAX/120 verified requests per - * second (one new id per ~120 ms) to evict legitimate entries. - * 1024 is generous for an admin endpoint. + * Default cap on the in-memory replay cache: with a 60s tolerance, 1024 live tokens is about + * 8 signed requests a second, generous for an admin endpoint. A public endpoint passes a + * larger `maxReplayEntries`. */ const val MAX_REPLAY_ENTRIES = 1024 + + /** Answered when the replay cache holds nothing but live tokens. */ + const val REPLAY_CACHE_FULL = "rate-limited: too many fresh NIP-98 tokens at once; retry shortly" } } From aef9f6f9a0a885bdb9c0b8063d496271c8a576ac Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 02:49:26 +0000 Subject: [PATCH 2/8] perf(quartz): keep NIP-77 snapshots per filter until a write can change their set LiveEventStore cached one sealed negentropy snapshot and dropped it on any accepted write. On a relay taking writes that meant a full scan + O(n log n) seal on nearly every NEG-OPEN, and NIP-FE's stateless /neg rounds re-open per round, so a reconcile paid it once per round instead of once per sync. The cache now keeps the last four filters (most recently used first) and drops an entry only when an accepted write can change its set: the event matches one of its filters; it is a replaceable or addressable event in the filter's kinds/authors (its predecessor may be in the set under tags or ids the new version lacks); or it is a deletion or vanish request (clears all). Deletes that bypass ingest stay bounded by the existing 30s TTL. The write generation counter is gone. LiveEventStoreSnapshotCacheTest pins each case against a counting store. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_016jDNhr6a3J4VC5TaG3Yd48 --- .../relay/server/backend/LiveEventStore.kt | 86 +++++++---- .../LiveEventStoreSnapshotCacheTest.kt | 133 ++++++++++++++++++ 2 files changed, 189 insertions(+), 30 deletions(-) create mode 100644 quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStoreSnapshotCacheTest.kt diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStore.kt index 2c37a5b904..1569cc194f 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStore.kt @@ -22,18 +22,21 @@ package com.vitorpamplona.quartz.nip01Core.relay.server.backend import com.vitorpamplona.negentropy.storage.IStorage import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.core.isAddressable +import com.vitorpamplona.quartz.nip01Core.core.isReplaceable import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.filters.FilterIndex import com.vitorpamplona.quartz.nip01Core.store.IEventStore import com.vitorpamplona.quartz.nip01Core.store.IdAndTime import com.vitorpamplona.quartz.nip01Core.store.RawEvent import com.vitorpamplona.quartz.nip01Core.store.StoreQueryContext +import com.vitorpamplona.quartz.nip09Deletions.DeletionRequestEvent +import com.vitorpamplona.quartz.nip62RequestToVanish.RequestToVanishEvent import com.vitorpamplona.quartz.utils.TimeUtils import kotlinx.coroutines.CompletableDeferred import kotlinx.coroutines.awaitCancellation import kotlinx.coroutines.withContext import kotlin.concurrent.atomics.AtomicBoolean -import kotlin.concurrent.atomics.AtomicLong import kotlin.concurrent.atomics.AtomicReference import kotlin.concurrent.atomics.ExperimentalAtomicApi @@ -154,7 +157,7 @@ class LiveEventStore( ) { ingest.submit(event, skipVerify) { outcome -> if (outcome is IEventStore.InsertOutcome.Accepted) { - writeGeneration.addAndFetch(1L) + forgetSnapshotsCovering(event) fanout(event) } onComplete(outcome) @@ -388,32 +391,28 @@ class LiveEventStore( // ------------------------------------------------------------------ /** - * Bumped after every accepted write. A cached negentropy snapshot is - * only valid while this hasn't moved. Deletion paths that bypass the - * ingest queue (expiration sweeps, NIP-86 admin purges) don't bump it, - * which is why cache entries also carry a short TTL: a snapshot is a - * point-in-time set by NIP-77's nature, and a few seconds of staleness - * only means a peer momentarily re-offers ids the relay just dropped. + * The last few sealed negentropy snapshots, most recently used first, each valid until a write + * could change its set. Replaced whole on every change, never mutated, so a reader never sees a + * half-updated list. Deletion paths that bypass the ingest queue (expiration sweeps, NIP-86 admin + * purges) are not seen here, which is why entries also carry a short TTL: a snapshot is a + * point-in-time set by NIP-77's nature, and a few seconds of staleness only means a peer + * momentarily re-offers ids the relay just dropped. */ - private val writeGeneration = AtomicLong(0L) + private val snapshotCache = AtomicReference>(emptyList()) private class CachedSnapshot( val filterKey: String, - val generation: Long, + val filters: List, val builtAt: Long, val storage: IStorage?, ) - private val snapshotCache = AtomicReference(null) - /** - * Serves repeated NEG-OPENs of the same filter from one sealed - * storage as long as no write landed in between (single slot — the - * mirror-heartbeat pattern is many peers reconciling the same broad - * filter, not many filters). Rebuilding on every open costs a full - * scan + O(n log n) seal that grows with the corpus: relayBench - * measured 342 ms per identical-set reconcile at 50k events vs - * strfry's 26 ms off its always-current tree. + * Serves repeated NEG-OPENs of the same filter from one sealed storage until a write lands in its + * set. Several slots, because NIP-FE's stateless rounds re-open per round and a relay syncs more + * than one filter at a time; rebuilding costs a full scan + O(n log n) seal that grows with the + * corpus: relayBench measured 342 ms per identical-set reconcile at 50k events vs strfry's 26 ms + * off its always-current tree. */ override suspend fun sealedNegentropyStorage( filters: List, @@ -428,30 +427,57 @@ class LiveEventStore( store.liveNegentropySnapshot(maxEntries)?.let { return it } } - val generation = writeGeneration.load() - val key = filters.joinToString("") { it.toJson() } + "cap=$maxEntries" + val key = filters.joinToString(" ") { it.toJson() } + " cap=$maxEntries" val now = TimeUtils.now() - val cached = snapshotCache.load() - if (cached != null && - cached.filterKey == key && - cached.generation == generation && - now - cached.builtAt <= SNAPSHOT_TTL_SECONDS - ) { - return cached.storage + snapshotCache.load().firstOrNull { it.filterKey == key && now - it.builtAt <= SNAPSHOT_TTL_SECONDS }?.let { hit -> + updateSnapshots { cached -> listOf(hit) + cached.filter { it !== hit } } + return hit.storage } val built = super.sealedNegentropyStorage(filters, maxEntries) - snapshotCache.store(CachedSnapshot(key, generation, now, built)) + val entry = CachedSnapshot(key, filters, now, built) + updateSnapshots { cached -> (listOf(entry) + cached.filter { it.filterKey != key }).take(SNAPSHOT_SLOTS) } return built } + /** + * Drops every cached snapshot whose set [event] can change: one of its filters admits the event, + * or the event can remove members it cannot see — a deletion or vanish request clears them all, + * and a replaceable or addressable event drops any snapshot whose kinds and authors it falls in, + * since the version it replaces may be in the set under tags or ids the new one no longer has. + */ + private fun forgetSnapshotsCovering(event: Event) { + if (snapshotCache.load().isEmpty()) return + if (event.kind == DeletionRequestEvent.KIND || event.kind == RequestToVanishEvent.KIND) { + updateSnapshots { emptyList() } + return + } + val supersedes = event.kind.isReplaceable() || event.kind.isAddressable() + updateSnapshots { cached -> + cached.filter { snapshot -> + snapshot.filters.none { f -> f.match(event) || (supersedes && Filter(kinds = f.kinds, authors = f.authors).match(event)) } + } + } + } + + private inline fun updateSnapshots(change: (List) -> List) { + while (true) { + val current = snapshotCache.load() + val next = change(current) + if (next == current || snapshotCache.compareAndSet(current, next)) return + } + } + private companion object { /** * Ceiling on how long a cached snapshot may serve NEG-OPENs even * with no observed writes — bounds staleness from delete paths - * the generation counter can't see. + * the ingest queue never sees. */ const val SNAPSHOT_TTL_SECONDS = 30L + + /** Distinct filters kept sealed at once; each holds up to maxSyncEvents ids at ~40 B each. */ + const val SNAPSHOT_SLOTS = 4 } } diff --git a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStoreSnapshotCacheTest.kt b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStoreSnapshotCacheTest.kt new file mode 100644 index 0000000000..eacc390d3f --- /dev/null +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStoreSnapshotCacheTest.kt @@ -0,0 +1,133 @@ +/* + * 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.Event +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync +import com.vitorpamplona.quartz.nip01Core.store.IEventStore +import com.vitorpamplona.quartz.nip01Core.store.IdAndTime +import com.vitorpamplona.quartz.nip01Core.store.sqlite.DefaultIndexingStrategy +import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore +import kotlinx.coroutines.CompletableDeferred +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.cancel +import kotlinx.coroutines.runBlocking +import kotlin.test.AfterTest +import kotlin.test.Test +import kotlin.test.assertEquals + +/** The NIP-77 snapshot cache: kept across writes that cannot change its set, dropped by those that can. */ +class LiveEventStoreSnapshotCacheTest { + /** Counts how often a snapshot is actually scanned out of the store. */ + private class Counting( + private val inner: IEventStore, + ) : IEventStore by inner { + var scans = 0 + + override suspend fun snapshotIdsForNegentropy( + filters: List, + maxEntries: Int?, + onProgress: ((collected: Int) -> Unit)?, + ): List { + scans++ + return inner.snapshotIdsForNegentropy(filters, maxEntries, onProgress) + } + } + + private val scope = CoroutineScope(Dispatchers.Default + SupervisorJob()) + private val store = Counting(EventStore(dbName = null, indexStrategy = DefaultIndexingStrategy(indexEventsByPubkeyAlone = true))) + private val live = LiveEventStore(store, IngestQueue(store, scope.coroutineContext)) + private val alice = NostrSignerSync() + private var clock = 1_700_000_000L + + private val notes = listOf(Filter(kinds = listOf(1))) + private val reactions = listOf(Filter(kinds = listOf(7))) + + @AfterTest + fun tearDown() { + scope.cancel() + } + + private fun event( + kind: Int, + tags: Array> = emptyArray(), + ) = alice.sign(clock++, kind, tags, "") + + private fun publish(event: Event) = + runBlocking { + val done = CompletableDeferred() + live.submit(event) { done.complete(it) } + assertEquals(IEventStore.InsertOutcome.Accepted, done.await()) + } + + private fun open(filters: List) = runBlocking { live.sealedNegentropyStorage(filters, maxEntries = 1_000) } + + @Test + fun twoFiltersInTurnAreEachScannedOnce() { + publish(event(1)) + repeat(3) { + open(notes) + open(reactions) + } + assertEquals(2, store.scans) + } + + @Test + fun aWriteOutsideTheSetKeepsTheSnapshot() { + open(notes) + publish(event(7)) + open(notes) + assertEquals(1, store.scans) + } + + @Test + fun aWriteInsideTheSetRebuildsIt() { + open(notes) + publish(event(1)) + open(notes) + assertEquals(2, store.scans) + } + + @Test + fun aDeletionDropsEverySnapshot() { + open(notes) + open(reactions) + publish(event(5, arrayOf(arrayOf("e", "a".repeat(64))))) + open(notes) + open(reactions) + assertEquals(4, store.scans) + } + + @Test + fun aReplacementDropsASnapshotItsOldVersionWasIn() { + // The stored profile is in the set by its tag; its replacement has no such tag and so does not + // match the filter, but it removes the old version from the set all the same. + val tagged = listOf(Filter(kinds = listOf(0), tags = mapOf("t" to listOf("nostr")))) + publish(event(0, arrayOf(arrayOf("t", "nostr")))) + open(tagged) + publish(event(0)) + open(tagged) + assertEquals(2, store.scans) + } +} From c2ae1578bca94abaad745937f97baeaf42777aa9 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 02:49:26 +0000 Subject: [PATCH 3/8] fix(quartz): NIP-FE review fixes, a policy vote for NIP-98 sign-in, frames without sub id MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Fixes from the review of #4212, each with a regression test in HttpRelayHandlerTest: - A deeply nested body threw StackOverflowError out of handle() (the kotlinx tree reader and re-serializer recurse). Bodies are no longer parsed into a tree: JsonShape checks the outer shape (one value, brackets matched by type, depth <= 32, nothing after it), the body is spliced into the frame, and the engine parses it once through RelaySession.receive(String). That also runs IRelayPolicy.acceptMessage again, which receive(Command) skipped, and drops the triple parse (tree, re-serialize, Jackson) to one. - A NIP-98 key was recorded without the policy chain's say, so a FullAuthPolicy whose authorize() refuses a user signed that user in over HTTP. New IRelayPolicy.acceptTransportIdentity(pubkey): the stack records the key only when no member objects, and the default objects. PassThroughPolicy and the verify policies have no objection; FullAuthPolicy asks a new authorizeTransport(pubkey), refusing by default, since authorize() has no AUTH event to run on. RelaySession.authenticateByTransport replaces the initialAuthenticatedUsers parameter #4212 added (and the matching connect/serve overloads), which recorded the key unconditionally. - The size gate compared UTF-8 bytes to a limit the engine counts in chars; it now measures the frame's chars, with a 3x byte bound for the read. - A backend or policy exception escaped handle(); it is now a 500 CLOSED line. - deadline is a Duration: Long.MAX_VALUE ms overflowed into an instant 503. - A reader stalled on a one-frame answer was never dropped; single() now has the same deadline-plus-grace bound as a stream. - An answer frame dequeued at the deadline was replaced by "ran past"; it now goes out, since the answer is complete. - The NIP-98 `u` is read with HTTPAuthorizationEvent.url(), as the verifier reads it; the status table maps MachineReadablePrefix exhaustively. - The handler's verifier is sized for a public endpoint (65,536 live tokens). And the wire format settled upstream in nostr-protocol/nips#2484: answers carry no subscription id — ["EVENT",{…}], ["EOSE"], ["COUNT",{…}], ["CLOSED",""]. The engine still runs under "http"; the handler strips it by verb, so frames that never carry one pass untouched. /neg stays as this implementation's extension (NIP-FE no longer names it): ["NEG-MSG",""]. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_016jDNhr6a3J4VC5TaG3Yd48 --- .../nip01Core/relay/server/RelayServerBase.kt | 21 +- .../nip01Core/relay/server/RelaySession.kt | 28 +- .../relay/server/policies/FullAuthPolicy.kt | 9 + .../relay/server/policies/IRelayPolicy.kt | 15 ++ .../server/policies/PassThroughPolicy.kt | 4 + .../relay/server/policies/PolicyStack.kt | 4 + .../relay/server/policies/VerifyPolicy.kt | 4 + .../nipFERelayOverHttp/HttpRelayCommand.kt | 200 +++++++++------ .../nipFERelayOverHttp/HttpRelayHandler.kt | 240 +++++++++++------- .../nipFERelayOverHttp/HttpRelayStatus.kt | 28 +- .../HttpRelayHandlerTest.kt | 188 ++++++++++++-- 11 files changed, 525 insertions(+), 216 deletions(-) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelayServerBase.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelayServerBase.kt index f1449f0261..b3fcb9ee5b 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelayServerBase.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelayServerBase.kt @@ -20,7 +20,6 @@ */ package com.vitorpamplona.quartz.nip01Core.relay.server -import com.vitorpamplona.quartz.nip01Core.core.HexKey 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.LimitsPolicy @@ -82,16 +81,8 @@ abstract class RelayServerBase( */ fun connect(send: (String) -> Unit): RelaySession = connect(SessionSink.of(send)) - /** - * Registers a new client connection whose frames go to [sink], typed where - * the engine has the type, already signed in as [authenticatedUsers] when - * the transport proved them itself (a NIP-98 header; see - * [RelaySession.initialAuthenticatedUsers]). - */ - fun connect( - sink: SessionSink, - authenticatedUsers: Set = emptySet(), - ): RelaySession = + /** Registers a new client connection whose frames go to [sink], typed where the engine has the type. */ + fun connect(sink: SessionSink): RelaySession = connections.register( RelaySession( policy = buildPolicy(), @@ -100,7 +91,6 @@ abstract class RelayServerBase( sink = sink, onClose = { connections.unregister(it.id) }, negentropySettings = negentropySettings, - initialAuthenticatedUsers = authenticatedUsers, ), ) @@ -115,15 +105,14 @@ abstract class RelayServerBase( suspend fun serve( send: (String) -> Unit, incoming: suspend (RelaySession) -> Unit, - ) = serve(SessionSink.of(send), emptySet(), incoming) + ) = serve(SessionSink.of(send), incoming) - /** [serve] over a [SessionSink], signed in as [authenticatedUsers] from the start. */ + /** [serve] over a [SessionSink]. */ suspend fun serve( sink: SessionSink, - authenticatedUsers: Set, incoming: suspend (RelaySession) -> Unit, ) { - val session = connect(sink, authenticatedUsers) + val session = connect(sink) try { incoming(session) } finally { diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelaySession.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelaySession.kt index 68793aa486..df77566329 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelaySession.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelaySession.kt @@ -79,14 +79,6 @@ class RelaySession( * open/close of the same connection. Defaults to a fresh monotonic id. */ val id: Long = nextConnectionId(), - /** - * Identities the transport proved before the session existed — a - * NIP-98 header on an HTTP command (NIP-FE), say — recorded exactly as - * a NIP-42 AUTH would record them. The policy's `accept(AuthCmd)` and - * `onAuthenticated` are not consulted: there is no AUTH event, and the - * transport, not the engine, did the verifying. - */ - initialAuthenticatedUsers: Set = emptySet(), ) : AutoCloseable { /** The original, string-only constructor; every frame goes to [onSend] as wire JSON. */ constructor( @@ -116,7 +108,7 @@ class RelaySession( * the copy costs nothing on the hot path. */ @Volatile - private var authenticatedUsers: Set = initialAuthenticatedUsers.toSet() + private var authenticatedUsers = setOf() /** * The per-connection scope. Handed to the [policy] at connect (so gating @@ -290,6 +282,24 @@ class RelaySession( send(CountMessage(cmd.queryId, countResult)) } + /** + * Records [pubkey], proved by the transport rather than a NIP-42 AUTH (NIP-FE's NIP-98 header), + * once the policy chain has no objection ([IRelayPolicy.acceptTransportIdentity]). Returns the + * refusal, or null when the key is now authenticated on this connection exactly as AUTH would. + */ + suspend fun authenticateByTransport(pubkey: HexKey): String? { + val refused = + try { + policy.acceptTransportIdentity(pubkey) + } catch (e: CancellationException) { + throw e + } catch (e: Exception) { + MachineReadablePrefix.ERROR.format(e.message ?: "authentication failed") + } + if (refused == null) authenticatedUsers = authenticatedUsers + pubkey + return refused + } + // -- NIP-42: AUTH --------------------------------------------------------- private suspend fun handleAuth(cmd: AuthCmd) { val result = policy.accept(cmd) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/FullAuthPolicy.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/FullAuthPolicy.kt index e13a6cddc8..814083271d 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/FullAuthPolicy.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/FullAuthPolicy.kt @@ -129,6 +129,15 @@ open class FullAuthPolicy( */ open suspend fun authorize(event: RelayAuthEvent) {} + final override suspend fun acceptTransportIdentity(pubkey: HexKey): String? = authorizeTransport(pubkey) + + /** + * [authorize]'s counterpart for a key the transport proved (NIP-98 on a NIP-FE command), which + * has no AUTH event to hand it. Refuses by default: a subclass whose [authorize] checks or grants + * anything must decide what that means without one, and opt in by returning null. + */ + open suspend fun authorizeTransport(pubkey: HexKey): String? = TRANSPORT_IDENTITY_REFUSED + override fun accept(cmd: EventCmd): PolicyResult = if (isAuthenticated()) { PolicyResult.Accepted(cmd) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/IRelayPolicy.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/IRelayPolicy.kt index 978512fee7..2c3440053b 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/IRelayPolicy.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/IRelayPolicy.kt @@ -21,6 +21,7 @@ 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 @@ -101,6 +102,17 @@ interface IRelayPolicy { */ suspend fun onAuthenticated(event: RelayAuthEvent): Boolean = false + /** + * This policy's say on recording [pubkey] as authenticated when the transport, not a NIP-42 + * AUTH, proved it: a NIP-98 header on a NIP-FE HTTP command. Return a NIP-01 reason to refuse, + * or null for no objection; the engine records the key only when no policy in the chain + * objects. There is no AUTH event here, so [accept] (AuthCmd) and [onAuthenticated] never run. + * + * The default objects, so a policy that decides logins in those two hooks is not bypassed by a + * transport it was not written for. Policies with no say over identity return null. + */ + suspend fun acceptTransportIdentity(pubkey: HexKey): String? = TRANSPORT_IDENTITY_REFUSED + /** * Inspects a raw inbound message before it is parsed. Return a reason * string to reject it (the engine sends it as a `NOTICE`), or null to let @@ -163,3 +175,6 @@ sealed interface PolicyResult { val reason: String, ) : PolicyResult } + +/** The refusal a policy that does not accept transport-proved identities answers with. */ +const val TRANSPORT_IDENTITY_REFUSED = "restricted: this relay signs in over NIP-42 AUTH only" diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/PassThroughPolicy.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/PassThroughPolicy.kt index 1d673f6e8b..109bd57908 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/PassThroughPolicy.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/PassThroughPolicy.kt @@ -21,6 +21,7 @@ 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.CountCmd @@ -52,4 +53,7 @@ open class PassThroughPolicy : IRelayPolicy { override fun accept(cmd: AuthCmd): PolicyResult = PolicyResult.Accepted(cmd) override fun canSendToSession(event: Event): Boolean = true + + /** No objection. A subclass that refuses logins in [accept] (AuthCmd) must refuse here too. */ + override suspend fun acceptTransportIdentity(pubkey: HexKey): String? = null } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/PolicyStack.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/PolicyStack.kt index c1062b874a..6adf69c413 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/PolicyStack.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/PolicyStack.kt @@ -21,6 +21,7 @@ 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 @@ -56,6 +57,9 @@ class PolicyStack( return policies.fold(false) { recorded, p -> p.onAuthenticated(event) || recorded } } + /** Every member must agree; the first objection is the answer. */ + override suspend fun acceptTransportIdentity(pubkey: HexKey): String? = policies.firstNotNullOfOrNull { it.acceptTransportIdentity(pubkey) } + override fun acceptMessage(message: String): String? = policies.firstNotNullOfOrNull { it.acceptMessage(message) } override fun acceptSubscription( diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/VerifyPolicy.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/VerifyPolicy.kt index a97c1431fa..9d7286d70b 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/VerifyPolicy.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/policies/VerifyPolicy.kt @@ -21,6 +21,7 @@ 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.crypto.verify import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.AuthCmd @@ -68,6 +69,9 @@ open class VerifyEventsAndAuthPolicy( } override fun canSendToSession(event: Event) = true + + /** No objection: this policy checks signatures, and the transport verified its own proof. */ + override suspend fun acceptTransportIdentity(pubkey: HexKey): String? = null } /** diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayCommand.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayCommand.kt index 304735bf06..6d48801955 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayCommand.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayCommand.kt @@ -20,35 +20,27 @@ */ package com.vitorpamplona.quartz.nipFERelayOverHttp -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 import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.EoseMessage import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.NoticeMessage import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.OkMessage -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.nip77Negentropy.NegErrMessage import com.vitorpamplona.quartz.nip77Negentropy.NegMsgMessage import com.vitorpamplona.quartz.nip77Negentropy.NegOpenCmd -import kotlinx.serialization.SerializationException -import kotlinx.serialization.json.Json -import kotlinx.serialization.json.JsonArray -import kotlinx.serialization.json.JsonElement -import kotlinx.serialization.json.JsonObject -import kotlinx.serialization.json.JsonPrimitive /** - * NIP-FE: the client commands HTTP carries, one path each. A body is the - * command's arguments after its subscription id (a lone object where the - * command takes one); the answer ends on the first frame [ends] accepts. + * NIP-FE: the client commands HTTP carries, one path each. A body is the command's arguments + * after its subscription id (a lone object where the command takes one); the answer ends on the + * first frame [ends] accepts. * - * `NEG` is one NIP-77 round, `[filter, message]`: the responder keeps no - * state between rounds but its snapshot, which the backend caches per - * filter, so each round carries its filter and there is no session to close. + * `NEG` is not in NIP-FE; it is this implementation's extension: one NIP-77 round, + * `[filter, message]`. The responder keeps no state between rounds but its snapshot, which the + * backend caches per filter, so each round carries its filter and there is no session to close. */ enum class HttpRelayCommand( val path: String, @@ -60,28 +52,35 @@ enum class HttpRelayCommand( ; /** - * The command [body] stands for, parsed and validated, or null when it - * is not this command's arguments. Re-serialized from a JSON tree, so a - * body can only ever be arguments, never a second command. + * The client frame [body] stands for, or null when it is not this command's arguments. Only + * the body's outer shape is checked here ([JsonShape]); the engine parses the frame once, so a + * malformed inside is its usual NOTICE. The shape check is what makes splicing safe: the body + * is one balanced value with nothing after it, so it cannot close the frame or open another. */ - fun parse(body: String): Parsed? { - val tree = - try { - Json.parseToJsonElement(body) - } catch (_: SerializationException) { - return null + fun frameOf(body: String): String? { + val shape = JsonShape.of(body) ?: return null + return when (this) { + REQ, COUNT -> { + val filters = + when { + shape.isObject -> shape.text + shape.elements.isNotEmpty() && shape.elements.all { it == '{' } -> shape.inner + else -> return null + } + frame(if (this == REQ) ReqCmd.LABEL else CountCmd.LABEL, SUB_ID, filters) } - val args = arguments(tree) ?: return null - val frame = JsonArray(head() + args).toString() - val cmd = runCatching { OptimizedJsonMapper.fromJsonToCommand(frame) }.getOrNull() ?: return null - return if (cmd.isValid() && cmd.matches()) Parsed(cmd, frame.length) else null - } - /** A body as its command, with the length of the frame it stands for: what the message-length limit measures. */ - class Parsed( - val command: Command, - val wireLength: Int, - ) + EVENT -> { + if (!shape.isObject) return null + "[\"${EventCmd.LABEL}\",${shape.text}]" + } + + NEG -> { + if (shape.isObject || shape.elements != NEG_ROUND) return null + frame(NegOpenCmd.LABEL, SUB_ID, shape.inner) + } + } + } /** Whether [message] is the last frame of this command's answer. */ fun ends(message: Message): Boolean = @@ -93,47 +92,106 @@ enum class HttpRelayCommand( NEG -> message is NegMsgMessage || message is NegErrMessage } - private fun head(): List = - when (this) { - REQ -> listOf(JsonPrimitive(ReqCmd.LABEL), JsonPrimitive(SUB_ID)) - COUNT -> listOf(JsonPrimitive(CountCmd.LABEL), JsonPrimitive(SUB_ID)) - EVENT -> listOf(JsonPrimitive(EventCmd.LABEL)) - NEG -> listOf(JsonPrimitive(NegOpenCmd.LABEL), JsonPrimitive(SUB_ID)) - } - - private fun arguments(body: JsonElement): List? = - when (this) { - REQ, COUNT -> { - when (body) { - is JsonObject -> listOf(body) - is JsonArray -> body.takeIf { it.isNotEmpty() && it.all { f -> f is JsonObject } } - else -> null - } - } - - EVENT -> { - (body as? JsonObject)?.let(::listOf) - } - - NEG -> { - (body as? JsonArray)?.takeIf { - it.size == 2 && it[0] is JsonObject && (it[1] as? JsonPrimitive)?.isString == true - } - } - } - - private fun Command.matches(): Boolean = - when (this@HttpRelayCommand) { - REQ -> this is ReqCmd - COUNT -> this is CountCmd - EVENT -> this is EventCmd - NEG -> this is NegOpenCmd - } - companion object { - /** The subscription id every HTTP command runs under; each request is its own connection. */ + /** + * The subscription id every HTTP command runs under inside the engine. NIP-FE answers carry + * none, so [HttpRelayHandler] takes it back out of each frame before it goes out. + */ const val SUB_ID = "http" + private val NEG_ROUND = listOf('{', '"') + fun forPath(path: String): HttpRelayCommand? = entries.firstOrNull { it.path == path } + + private fun frame( + label: String, + subId: String, + args: String, + ) = "[\"$label\",\"$subId\",$args]" + } +} + +/** + * The outer shape of a JSON body, read without building a tree: one object or array, brackets + * matched by type outside strings, nesting no deeper than [MAX_DEPTH], nothing after it. An array's + * [elements] are each top-level element's first character. Everything inside is left to the parser. + */ +internal class JsonShape private constructor( + val text: String, + val isObject: Boolean, + val elements: List, +) { + /** An array's contents without its brackets. */ + val inner: String get() = text.substring(1, text.length - 1) + + companion object { + /** Deep enough for any filter, event or round by a wide margin; far too shallow to exhaust a stack. */ + const val MAX_DEPTH = 32 + + fun of(body: String): JsonShape? { + val text = body.trim() + if (text.length < 2 || (text[0] != '{' && text[0] != '[')) return null + val elements = ArrayList() + val open = CharArray(MAX_DEPTH) + var depth = 0 + var inString = false + var escaped = false + // At depth 1 inside an array: whether the next non-space character starts an element. + var expectElement = text[0] == '[' + var i = 0 + while (i < text.length) { + val c = text[i] + if (inString) { + when { + escaped -> escaped = false + c == '\\' -> escaped = true + c == '"' -> inString = false + } + i++ + continue + } + if (depth == 0 && i > 0) return null + if (depth == 1 && text[0] == '[' && !c.isWhitespace()) { + when { + c == ',' -> { + if (expectElement) return null + expectElement = true + i++ + continue + } + + c == ']' -> { + if (expectElement && elements.isNotEmpty()) return null + } + + expectElement -> { + elements.add(c) + expectElement = false + } + + c == '{' || c == '[' || c == '"' -> { + return null + } + } + } + when (c) { + '"' -> { + inString = true + } + + '{', '[' -> { + if (depth == MAX_DEPTH) return null + open[depth++] = c + } + + '}', ']' -> { + if (depth == 0 || open[--depth] != (if (c == '}') '{' else '[')) return null + } + } + i++ + } + if (depth != 0 || inString) return null + return JsonShape(text, text[0] == '{', elements) + } } } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt index e28f86815d..e443028888 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt @@ -28,7 +28,9 @@ import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.MachineReadabl import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message import com.vitorpamplona.quartz.nip01Core.relay.server.RelayServerBase import com.vitorpamplona.quartz.nip01Core.relay.server.SessionSink +import com.vitorpamplona.quartz.nip98HttpAuth.HTTPAuthorizationEvent import com.vitorpamplona.quartz.nip98HttpAuth.Nip98AuthVerifier +import kotlinx.coroutines.CancellationException import kotlinx.coroutines.CompletableDeferred import kotlinx.coroutines.TimeoutCancellationException import kotlinx.coroutines.channels.Channel @@ -38,20 +40,27 @@ import kotlinx.coroutines.withTimeout import kotlinx.coroutines.withTimeoutOrNull import kotlin.io.encoding.Base64 import kotlin.io.encoding.ExperimentalEncodingApi +import kotlin.time.Duration +import kotlin.time.Duration.Companion.milliseconds import kotlin.time.TimeSource -/** One NIP-FE request as the handler needs it. The host routes by [HttpRelayCommand.path] and bounds the body while reading it. */ +/** One NIP-FE request as the handler needs it. The host routes by [HttpRelayCommand.path]. */ class HttpRelayRequest( val command: HttpRelayCommand, /** The `Authorization` header as sent, or null. */ val authorization: String?, + /** + * The body. Hosts bound the read at [HttpRelayHandler.maxBodyBytes]: the engine measures a + * frame in characters, and a UTF-8 character takes up to three bytes. + */ val body: ByteArray, ) /** - * Where [HttpRelayHandler] writes an answer. The host owns the socket, the - * headers and any compression; the handler decides the status and the lines. - * On 401 the host adds `WWW-Authenticate: Nostr`, on 429 and 503 `Retry-After`. + * Where [HttpRelayHandler] writes an answer. The host owns the socket, the headers and any + * compression; the handler decides the status and the lines. On 401 the host adds + * `WWW-Authenticate: Nostr`, on 429 and 503 `Retry-After`. A [HttpRelayReaderStalled] thrown out of + * either call means the client stopped reading, and the host drops the connection unfinished. */ interface HttpRelayResponse { /** An answer that is one frame, with its status: a refusal, or a command answered at once. */ @@ -60,11 +69,7 @@ interface HttpRelayResponse { frame: String, ) - /** - * A 200 answer written line by line, as `application/x-ndjson`. Returns - * when [lines] does; a [HttpRelayReaderStalled] thrown out of it means - * the client stopped reading and the host drops the connection unfinished. - */ + /** A 200 answer written line by line, as `application/x-ndjson`. Returns when [lines] does. */ suspend fun stream(lines: suspend HttpRelayLines.() -> Unit) } @@ -77,129 +82,148 @@ interface HttpRelayLines { suspend fun flush() } -/** The client stopped reading a streamed answer; the host drops the connection instead of finishing it. */ +/** The client stopped reading an answer; the host drops the connection instead of finishing it. */ class HttpRelayReaderStalled : Exception("the client stopped reading the answer") /** - * NIP-FE: one relay command per HTTP request, run on its own [RelayServerBase] - * session, so every limit and policy the socket applies applies here, and - * answered with the relay's own frames up to the command's answer. Nothing - * outlives the request. Admission (how many requests a client may run) is the - * host's: gate before calling [handle], so a refused request does not spend a - * NIP-98 token that the handler would have verified. + * NIP-FE: one relay command per HTTP request, run on its own [RelayServerBase] session, so every + * limit and policy the socket applies applies here, and answered with the relay's own frames up to + * the command's answer, without their subscription id. Nothing outlives the request. + * + * Admission (how many requests a client may run) is the host's: gate before calling [handle], so a + * refused request does not spend a NIP-98 token the handler would have verified. */ class HttpRelayHandler( private val server: RelayServerBase, /** The prefixes a NIP-98 `u` may carry (the relay's http origin, its .onion), asked per request; never from the request. */ private val origins: () -> List, - /** How long one answer may run, first byte to last. */ - private val deadlineMs: Long = DEFAULT_DEADLINE_MS, - /** Its own replay cache, so public commands cannot evict another endpoint's. */ - private val verifier: Nip98AuthVerifier = Nip98AuthVerifier(), + /** How long one answer may run, first byte to last. [Duration.INFINITE] turns the deadline off. */ + private val deadline: Duration = DEFAULT_DEADLINE, + /** Its own replay cache, sized for a public endpoint, so it cannot be flushed to replay a token. */ + private val verifier: Nip98AuthVerifier = Nip98AuthVerifier(maxReplayEntries = DEFAULT_REPLAY_ENTRIES), /** Frames queued ahead of a slow reader before the answer is cut short. */ private val maxQueuedFrames: Int = DEFAULT_MAX_QUEUED_FRAMES, /** How long past the deadline the last line may take before the reader counts as stalled. */ - private val tailGraceMs: Long = DEFAULT_TAIL_GRACE_MS, + private val tailGrace: Duration = DEFAULT_TAIL_GRACE, ) { + /** The largest body that can still be a frame within the relay's message limit, or null for no limit. */ + val maxBodyBytes: Long? get() = server.limits?.maxMessageLength?.let { it.toLong() * 3 } + suspend fun handle( request: HttpRelayRequest, response: HttpRelayResponse, ) { val command = request.command val max = server.limits?.maxMessageLength - if (max != null && request.body.size > max) { - return response.single(HttpRelayStatus.PAYLOAD_TOO_LARGE, closed("invalid: the body exceeds $max bytes")) + maxBodyBytes?.let { cap -> + if (request.body.size > cap) { + return response.single(HttpRelayStatus.PAYLOAD_TOO_LARGE, closed("invalid: the command exceeds $max characters")) + } } - val parsed = - command.parse(request.body.decodeToString()) + val frame = + command.frameOf(request.body.decodeToString()) ?: return response.single(HttpRelayStatus.BAD_REQUEST, closed("invalid: the body is not ${command.name}'s arguments")) - if (max != null && parsed.wireLength > max) { + // Characters, as the engine's own limit counts them. + if (max != null && frame.length > max) { return response.single(HttpRelayStatus.PAYLOAD_TOO_LARGE, closed("invalid: the command exceeds $max characters")) } val signedIn = when (val proof = proofOf(request)) { - is Proof.Anonymous -> emptySet() - is Proof.Signed -> setOf(proof.pubkey) - is Proof.Refused -> return response.single(HttpRelayStatus.UNAUTHORIZED, closed(MachineReadablePrefix.AUTH_REQUIRED.format(proof.reason))) + is Proof.Anonymous -> { + null + } + + is Proof.Signed -> { + proof.pubkey + } + + is Proof.Refused -> { + val reason = proof.reason + return response.single(HttpRelayStatus.forReason(reason), closed(reason)) + } } - exchange(parsed, command, signedIn, response) + exchange(frame, command, signedIn, response) } - /** A frame as queued: its wire text, and its type when the engine built one. */ + /** A frame as queued: its wire text, its type when the engine built one, and whether it ends the answer. */ private class Frame( val json: String, val message: Message?, + val last: Boolean, ) - private fun HttpRelayCommand.ends(frame: Frame) = frame.message?.let(::ends) == true - private suspend fun exchange( - parsed: HttpRelayCommand.Parsed, + frame: String, command: HttpRelayCommand, - signedIn: Set, + signedIn: HexKey?, response: HttpRelayResponse, ) = coroutineScope { val frames = Channel(maxQueuedFrames) val ended = CompletableDeferred() // Called on the engine's coroutines and cannot suspend, so a frame that does not fit ends the answer. - fun offer( - frame: Frame, - last: Boolean, - ) { + fun offer(frame: Frame) { val sent = frames.trySend(frame) if (sent.isClosed) return - if (sent.isFailure || last) { + if (sent.isFailure || frame.last) { frames.close() ended.complete(Unit) } } + + fun fail(reason: String) = offer(Frame(closed(reason), ClosedMessage(HttpRelayCommand.SUB_ID, reason), last = true)) val sink = object : SessionSink { override fun message(message: Message) { - // The challenge every connection opens with; this one proved its key by NIP-98 instead. + // The challenge every connection opens with; this one proves its key by NIP-98 instead. if (message is AuthMessage) return - offer(Frame(message.toJson(), message), command.ends(message)) + offer(Frame(withoutSubId(message.toJson()), message, command.ends(message))) } - override fun raw(json: String) = offer(Frame(json, null), false) + override fun raw(json: String) = offer(Frame(withoutSubId(json), null, last = false)) } val session = launch { - server.serve(sink, signedIn) { - it.receive(parsed.command) - ended.await() + try { + server.serve(sink) { session -> + val refused = signedIn?.let { session.authenticateByTransport(it) } + if (refused != null) fail(refused) else session.receive(frame) + ended.await() + } + } catch (e: CancellationException) { + throw e + } catch (e: Exception) { + // A backend or policy that throws is the relay's failure, answered as one. + fail(MachineReadablePrefix.ERROR.format(e.message ?: "the relay failed this command")) } } - val started = TimeSource.Monotonic.markNow() + val due = TimeSource.Monotonic.markNow() + deadline - fun remainingMs() = (deadlineMs - started.elapsedNow().inWholeMilliseconds).coerceAtLeast(0) + suspend fun single( + status: Int, + json: String, + ) = bounded(due) { response.single(status, json) } try { - val first = withTimeoutOrNull(deadlineMs) { frames.receiveCatching().getOrNull() } + val first = withTimeoutOrNull(deadline) { frames.receiveCatching().getOrNull() } + val status = HttpRelayStatus.of(first?.message) when { first == null -> { - response.single(HttpRelayStatus.UNAVAILABLE, closed("error: no answer within ${deadlineMs / 1000}s")) + single(HttpRelayStatus.UNAVAILABLE, closed("error: no answer within $deadline")) } - command.ends(first) || HttpRelayStatus.of(first.message) != HttpRelayStatus.OK -> { - response.single(HttpRelayStatus.of(first.message), first.json) + first.last || status != HttpRelayStatus.OK -> { + single(status, first.json) } else -> { response.stream { - // The deadline is read between frames and never interrupts a write, so every line leaves - // whole; a reader that stops reading altogether is dropped at the hard stop. - try { - withTimeout(remainingMs() + tailGraceMs) { - when (drain(first, frames, command, started)) { - Ending.ANSWERED -> {} - Ending.DEADLINE -> line(closed("error: the answer ran past ${deadlineMs / 1000}s")) - Ending.CUT -> line(closed("error: slow reader, over $maxQueuedFrames frames waiting")) - } - flush() + bounded(due) { + when (drain(first, frames, due)) { + Ending.ANSWERED -> {} + Ending.DEADLINE -> line(closed("error: the answer ran past $deadline")) + Ending.CUT -> line(closed("error: slow reader, over $maxQueuedFrames frames waiting")) } - } catch (_: TimeoutCancellationException) { - throw HttpRelayReaderStalled() + flush() } } } @@ -209,32 +233,46 @@ class HttpRelayHandler( } } + /** Runs [block] until [due] plus the tail grace; a write still blocked then is a reader that stopped. */ + private suspend fun bounded( + due: TimeSource.Monotonic.ValueTimeMark, + block: suspend () -> Unit, + ) { + val left = (-due.elapsedNow()).coerceAtLeast(Duration.ZERO) + tailGrace + try { + withTimeout(left) { block() } + } catch (_: TimeoutCancellationException) { + throw HttpRelayReaderStalled() + } + } + /** How a streamed answer stopped: at its answer frame, at the deadline, or cut because the reader fell behind. */ private enum class Ending { ANSWERED, DEADLINE, CUT } /** - * Writes [first] and what follows up to the command's answer, flushing whenever nothing is waiting so - * a burst leaves as one write. Stops at the answer: a live event queued behind it is not part of it. + * Writes [first] and what follows up to the command's answer, flushing whenever nothing is + * waiting so a burst leaves as one write. Stops at the answer: a live event queued behind it is + * not part of it. The deadline is read between frames and never interrupts a write, so every line + * leaves whole, and the answer frame goes out even at the deadline: the answer is complete. */ private suspend fun HttpRelayLines.drain( first: Frame, frames: Channel, - command: HttpRelayCommand, - started: TimeSource.Monotonic.ValueTimeMark, + due: TimeSource.Monotonic.ValueTimeMark, ): Ending { var frame = first while (true) { line(frame.json) - if (command.ends(frame)) return Ending.ANSWERED + if (frame.last) return Ending.ANSWERED frame = frames.tryReceive().getOrNull() ?: run { flush() - val leftMs = deadlineMs - started.elapsedNow().inWholeMilliseconds - if (leftMs <= 0) return Ending.DEADLINE - val next = withTimeoutOrNull(leftMs) { frames.receiveCatching() } ?: return Ending.DEADLINE + val left = -due.elapsedNow() + if (!left.isPositive()) return Ending.DEADLINE + val next = withTimeoutOrNull(left) { frames.receiveCatching() } ?: return Ending.DEADLINE // Closed with no answer frame in it: the send side gave up on this reader. next.getOrNull() ?: return Ending.CUT } - if (started.elapsedNow().inWholeMilliseconds >= deadlineMs) return Ending.DEADLINE + if (!frame.last && due.hasPassedNow()) return Ending.DEADLINE } } @@ -262,34 +300,64 @@ class HttpRelayHandler( if (!header.regionMatches(0, scheme, 0, scheme.length, ignoreCase = true)) return Proof.Anonymous val token = scheme + header.substring(scheme.length).trim() val accepted = origins().map { it.trimEnd('/') + request.command.path } - val url = claimedUrl(token)?.takeIf { it in accepted } ?: accepted.firstOrNull() ?: return Proof.Refused("this relay names no url to sign") + val url = claimedUrl(token)?.takeIf { it in accepted } ?: accepted.firstOrNull() ?: return Proof.Refused(MachineReadablePrefix.AUTH_REQUIRED.format("this relay names no url to sign")) return when (val r = verifier.verify(token, "POST", url, request.body)) { - is Nip98AuthVerifier.Result.Verified -> Proof.Signed(r.pubkey) - is Nip98AuthVerifier.Result.Malformed -> Proof.Refused("NIP-98 ${r.reason}") - is Nip98AuthVerifier.Result.Missing -> Proof.Anonymous + is Nip98AuthVerifier.Result.Verified -> { + Proof.Signed(r.pubkey) + } + + is Nip98AuthVerifier.Result.Missing -> { + Proof.Anonymous + } + + // A full replay cache is the relay's limit, not the token's fault. + is Nip98AuthVerifier.Result.Malformed -> { + if (MachineReadablePrefix.parse(r.reason) == MachineReadablePrefix.RATE_LIMITED) { + Proof.Refused(r.reason) + } else { + Proof.Refused(MachineReadablePrefix.AUTH_REQUIRED.format("NIP-98 ${r.reason}")) + } + } } } - /** The `u` tag of a NIP-98 token, or null when it does not decode; the verifier then says why. */ + /** The `u` a NIP-98 token names, read as the verifier reads it, or null when it does not decode. */ @OptIn(ExperimentalEncodingApi::class) private fun claimedUrl(token: String): String? = runCatching { val json = Base64.decode(token.removePrefix(Nip98AuthVerifier.SCHEME).trim()).decodeToString() - OptimizedJsonMapper - .fromJson(json) - .tags - .firstOrNull { it.size > 1 && it[0] == "u" } - ?.get(1) + val event = OptimizedJsonMapper.fromJson(json) + HTTPAuthorizationEvent(event.id, event.pubKey, event.createdAt, event.tags, event.content, event.sig).url() }.getOrNull() - private fun closed(reason: String) = ClosedMessage(HttpRelayCommand.SUB_ID, reason).toJson() + private fun closed(reason: String) = withoutSubId(ClosedMessage(HttpRelayCommand.SUB_ID, reason).toJson()) companion object { - const val DEFAULT_DEADLINE_MS = 30_000L + val DEFAULT_DEADLINE = 30_000.milliseconds /** The websocket's slow-consumer bound in the reference relays. */ const val DEFAULT_MAX_QUEUED_FRAMES = 8192 - const val DEFAULT_TAIL_GRACE_MS = 5_000L + val DEFAULT_TAIL_GRACE = 5_000.milliseconds + + /** Two minutes of tokens (the replay window) at about 500 signed commands a second. */ + const val DEFAULT_REPLAY_ENTRIES = 65_536 } } + +/** The frames that carry a subscription id in the engine; NIP-FE sends them without it. */ +private val SUBSCRIPTION_FRAMES = setOf("EVENT", "EOSE", "CLOSED", "COUNT", "NEG-MSG", "NEG-ERR") + +private const val SUB_ID_FIELD = ",\"" + HttpRelayCommand.SUB_ID + "\"" + +/** + * [frame] as NIP-FE sends it: the engine's frame with its `"http"` subscription id taken out, + * `["EVENT","http",{…}]` → `["EVENT",{…}]`, `["EOSE","http"]` → `["EOSE"]`. Other frames pass as they are. + */ +internal fun withoutSubId(frame: String): String { + if (!frame.startsWith("[\"")) return frame + val verbEnd = frame.indexOf('"', 2) + if (verbEnd < 0 || frame.substring(2, verbEnd) !in SUBSCRIPTION_FRAMES) return frame + if (!frame.startsWith(SUB_ID_FIELD, verbEnd + 1)) return frame + return frame.substring(0, verbEnd + 1) + frame.substring(verbEnd + 1 + SUB_ID_FIELD.length) +} diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayStatus.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayStatus.kt index 41f8fb64a2..496acab463 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayStatus.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayStatus.kt @@ -21,16 +21,17 @@ package com.vitorpamplona.quartz.nipFERelayOverHttp import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.ClosedMessage +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.NoticeMessage import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.OkMessage +import com.vitorpamplona.quartz.nip01Core.store.RejectionReason import com.vitorpamplona.quartz.nip77Negentropy.NegErrMessage /** - * NIP-FE status codes. The status of an answer is decided by its first - * frame: an accepting one is 200 and the answer may stream; a refusal's - * NIP-01 machine-readable prefix picks the code. After the first frame is - * out the status cannot change, so later failures are frames. + * NIP-FE status codes. The status of an answer is decided by its first frame: an accepting one + * is 200 and the answer may stream; a refusal's NIP-01 machine-readable prefix picks the code. + * After the first frame is out the status cannot change, so later failures are frames. */ object HttpRelayStatus { const val OK = 200 @@ -42,24 +43,25 @@ object HttpRelayStatus { const val INTERNAL_ERROR = 500 const val UNAVAILABLE = 503 - /** The status an answer opening with [message] carries. A raw EVENT frame opens a 200. */ + /** The status an answer opening with [message] carries. A raw EVENT frame (no type) opens a 200. */ fun of(message: Message?): Int = when (message) { is ClosedMessage -> forReason(message.message) is NegErrMessage -> forReason(message.reason) // A duplicate is already stored, which is what the caller asked for, whichever flag the store set. - is OkMessage -> if (message.success || message.message.startsWith("duplicate:")) OK else forReason(message.message) + is OkMessage -> if (message.success || message.message.startsWith(RejectionReason.PREFIX_DUPLICATE)) OK else forReason(message.message) is NoticeMessage -> BAD_REQUEST else -> OK } - /** The status a NIP-01 machine-readable prefix stands for. */ + /** The status a NIP-01 machine-readable prefix stands for; a reason without one is a 400. */ fun forReason(reason: String): Int = - when (reason.substringBefore(':', "")) { - "auth-required" -> UNAUTHORIZED - "restricted", "blocked" -> FORBIDDEN - "rate-limited" -> TOO_MANY_REQUESTS - "error" -> INTERNAL_ERROR - else -> BAD_REQUEST + when (MachineReadablePrefix.parse(reason)) { + MachineReadablePrefix.AUTH_REQUIRED -> UNAUTHORIZED + MachineReadablePrefix.RESTRICTED, MachineReadablePrefix.BLOCKED -> FORBIDDEN + MachineReadablePrefix.RATE_LIMITED -> TOO_MANY_REQUESTS + MachineReadablePrefix.ERROR -> INTERNAL_ERROR + MachineReadablePrefix.DUPLICATE -> OK + MachineReadablePrefix.INVALID, MachineReadablePrefix.POW, MachineReadablePrefix.UNSUPPORTED, null -> BAD_REQUEST } } diff --git a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt index e6837eb856..4123405322 100644 --- a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt @@ -23,15 +23,20 @@ package com.vitorpamplona.quartz.nipFERelayOverHttp import com.vitorpamplona.negentropy.Negentropy import com.vitorpamplona.negentropy.storage.StorageVector import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.core.HexKey import com.vitorpamplona.quartz.nip01Core.core.toHexKey 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.toRelay.ReqCmd import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer import com.vitorpamplona.quartz.nip01Core.relay.server.RelayServerBase import com.vitorpamplona.quartz.nip01Core.relay.server.RelayServerListener 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.FullAuthPolicy +import com.vitorpamplona.quartz.nip01Core.relay.server.policies.IRelayPolicy +import com.vitorpamplona.quartz.nip01Core.relay.server.policies.LimitsPolicy import com.vitorpamplona.quartz.nip01Core.relay.server.policies.PassThroughPolicy import com.vitorpamplona.quartz.nip01Core.relay.server.policies.PolicyResult import com.vitorpamplona.quartz.nip01Core.relay.server.policies.RelayLimits @@ -39,8 +44,10 @@ import com.vitorpamplona.quartz.nip01Core.relay.server.policies.VerifyPolicy import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync import com.vitorpamplona.quartz.nip01Core.store.IEventStore import com.vitorpamplona.quartz.nip01Core.store.IdAndTime +import com.vitorpamplona.quartz.nip42RelayAuth.RelayAuthEvent import com.vitorpamplona.quartz.nip77Negentropy.NegentropySettings import com.vitorpamplona.quartz.nip98HttpAuth.HTTPAuthorizationEvent +import com.vitorpamplona.quartz.nip98HttpAuth.Nip98AuthVerifier import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.awaitCancellation import kotlinx.coroutines.runBlocking @@ -49,6 +56,9 @@ import kotlin.test.Test import kotlin.test.assertEquals import kotlin.test.assertFailsWith import kotlin.test.assertTrue +import kotlin.time.Duration +import kotlin.time.Duration.Companion.milliseconds +import kotlin.time.Duration.Companion.seconds /** NIP-FE's handler over a relay engine and an in-memory backend: status, lines, auth, and negentropy rounds. */ class HttpRelayHandlerTest { @@ -108,15 +118,19 @@ class HttpRelayHandlerTest { } private class MemoryRelay( - override val backend: MemoryBackend, - signedInOnly: Boolean, + override val backend: SessionBackend, + policies: () -> IRelayPolicy, + limits: RelayLimits? = RelayLimits(maxMessageLength = 4096), ) : RelayServerBase( - policyBuilder = { if (signedInOnly) VerifyPolicy + SignedInOnly() else VerifyPolicy }, + policyBuilder = policies, parentContext = SupervisorJob(), negentropySettings = NegentropySettings.Default, listener = RelayServerListener.None, - limits = RelayLimits(maxMessageLength = 4096), - ) + limits = limits, + ) { + constructor(backend: SessionBackend, signedInOnly: Boolean) : + this(backend, { if (signedInOnly) VerifyPolicy + SignedInOnly() else VerifyPolicy }) + } /** One answer as the host would see it. */ private class Recorded : HttpRelayResponse { @@ -152,8 +166,8 @@ class HttpRelayHandlerTest { private fun handler( signedInOnly: Boolean = false, - deadlineMs: Long = 5_000, - ) = HttpRelayHandler(MemoryRelay(backend, signedInOnly), origins = { listOf(origin) }, deadlineMs = deadlineMs) + deadline: Duration = 5.seconds, + ) = HttpRelayHandler(MemoryRelay(backend, signedInOnly), origins = { listOf(origin) }, deadline = deadline) private fun HttpRelayHandler.ask( command: HttpRelayCommand, @@ -179,7 +193,7 @@ class HttpRelayHandlerTest { val answer = handler().ask(HttpRelayCommand.REQ, """{"kinds":[1]}""") assertEquals(200, answer.status) assertTrue(answer.streamed) - assertEquals("""["EOSE","http"]""", answer.lines.last()) + assertEquals("""["EOSE"]""", answer.lines.last()) assertEquals( setOf(a.id, b.id), answer.lines @@ -193,7 +207,7 @@ class HttpRelayHandlerTest { fun anEmptyReqIsOneEoseLine() { val answer = handler().ask(HttpRelayCommand.REQ, """[{"kinds":[30000]}]""") assertEquals(200, answer.status) - assertEquals(listOf("""["EOSE","http"]"""), answer.lines) + assertEquals(listOf("""["EOSE"]"""), answer.lines) } @Test @@ -215,22 +229,26 @@ class HttpRelayHandlerTest { backend.events += listOf(note("a"), note("b"), note("c")) val answer = handler().ask(HttpRelayCommand.COUNT, """{"kinds":[1]}""") assertEquals(200, answer.status) - assertTrue(answer.lines.single().startsWith("""["COUNT","http",{"count":3"""), answer.lines.toString()) + assertTrue(answer.lines.single().startsWith("""["COUNT",{"count":3"""), answer.lines.toString()) } @Test fun aBodyThatIsNotTheCommandsArgumentsIsA400AndOneOverTheLimitA413() { + // Not the command's shape: refused before any session opens. for ((command, body) in listOf( HttpRelayCommand.REQ to "[]", HttpRelayCommand.EVENT to """[{"id":"x"}]""", - HttpRelayCommand.EVENT to """{"id":"not an event"}""", HttpRelayCommand.NEG to """[{"kinds":[1]}]""", HttpRelayCommand.COUNT to "not json", )) { val answer = handler().ask(command, body) assertEquals(400, answer.status, "$command '$body'") - assertTrue(answer.lines.single().startsWith("""["CLOSED","http","invalid:"""), answer.lines.toString()) + assertTrue(answer.lines.single().startsWith("""["CLOSED","invalid:"""), answer.lines.toString()) } + // The right shape with an inside the engine cannot read: its own NOTICE, the command never ran. + val unreadable = handler().ask(HttpRelayCommand.EVENT, """{"id":"not an event"}""") + assertEquals(400, unreadable.status) + assertTrue(unreadable.lines.single().startsWith("""["NOTICE","""), unreadable.lines.toString()) assertEquals(413, handler().ask(HttpRelayCommand.REQ, """{"search":"${"x".repeat(5_000)}"}""").status) // Under the byte cap, over it once wrapped in its frame: the engine measures the frame. assertEquals(413, handler().ask(HttpRelayCommand.REQ, """{"search":"${"x".repeat(4_096 - 20)}"}""").status) @@ -244,13 +262,13 @@ class HttpRelayHandlerTest { val anonymous = gated.ask(HttpRelayCommand.REQ, body) assertEquals(401, anonymous.status) - assertTrue(anonymous.lines.single().startsWith("""["CLOSED","http","auth-required:""")) + assertTrue(anonymous.lines.single().startsWith("""["CLOSED","auth-required:""")) assertEquals(401, gated.ask(HttpRelayCommand.REQ, body, "Basic dXNlcjpwYXNz").status, "Basic is not addressed to the relay") val signed = gated.ask(HttpRelayCommand.REQ, body, token(HttpRelayCommand.REQ, body)) assertEquals(200, signed.status, signed.lines.toString()) - assertEquals("""["EOSE","http"]""", signed.lines.last()) + assertEquals("""["EOSE"]""", signed.lines.last()) } @Test @@ -288,7 +306,7 @@ class HttpRelayHandlerTest { val reply = answer.lines .single() - .substringAfter("""["NEG-MSG","http","""") + .substringAfter("""["NEG-MSG","""") .substringBefore('"') val result = client.reconcile(reply.hexToByteArray()) have += result.sendIds.map { it.toHexString() } @@ -303,30 +321,30 @@ class HttpRelayHandlerTest { fun aMalformedNegentropyRoundIsRefused() { val answer = handler().ask(HttpRelayCommand.NEG, """[{"kinds":[1]},"zz"]""") assertEquals(400, answer.status) - assertTrue(answer.lines.single().startsWith("""["NEG-ERR","http","""), answer.lines.toString()) + assertTrue(answer.lines.single().startsWith("""["NEG-ERR","""), answer.lines.toString()) } @Test fun noFirstFrameWithinTheDeadlineIsA503() { - val answer = handler(deadlineMs = 300).ask(HttpRelayCommand.REQ, """{"kinds":[$STALLED_KIND]}""") + val answer = handler(deadline = 300.milliseconds).ask(HttpRelayCommand.REQ, """{"kinds":[$STALLED_KIND]}""") assertEquals(503, answer.status) - assertTrue(answer.lines.single().startsWith("""["CLOSED","http","error: no answer""")) + assertTrue(answer.lines.single().startsWith("""["CLOSED","error: no answer""")) } @Test fun aDeadlineMidAnswerEndsOnAClosedLine() { val found = note("found") backend.events += found - val answer = handler(deadlineMs = 300).ask(HttpRelayCommand.REQ, """{"kinds":[1,$TRICKLE_KIND]}""") + val answer = handler(deadline = 300.milliseconds).ask(HttpRelayCommand.REQ, """{"kinds":[1,$TRICKLE_KIND]}""") assertEquals(200, answer.status) assertTrue(found.id in answer.lines.first()) - assertTrue(answer.lines.last().startsWith("""["CLOSED","http","error: the answer ran past"""), answer.lines.toString()) + assertTrue(answer.lines.last().startsWith("""["CLOSED","error: the answer ran past"""), answer.lines.toString()) } @Test fun aReaderThatStopsReadingIsDroppedAtTheHardStop() { backend.events += note("found") - val h = HttpRelayHandler(MemoryRelay(backend, false), origins = { listOf(origin) }, deadlineMs = 100, tailGraceMs = 100) + val h = HttpRelayHandler(MemoryRelay(backend, false), origins = { listOf(origin) }, deadline = 100.milliseconds, tailGrace = 100.milliseconds) val stalled = object : HttpRelayResponse { override suspend fun single( @@ -356,6 +374,134 @@ class HttpRelayHandlerTest { session.close() } + @Test + fun framesCarryNoSubscriptionId() { + val found = note("found") + backend.events += found + val answer = handler().ask(HttpRelayCommand.REQ, """{"kinds":[1]}""") + assertTrue(answer.lines.first().startsWith("""["EVENT",{"""), answer.lines.toString()) + assertEquals("""["EOSE"]""", answer.lines.last()) + } + + @Test + fun aDeeplyNestedBodyIsA400NotAStackOverflow() { + for (body in listOf("""{"a":""".repeat(2_000) + "1" + "}".repeat(2_000), "[".repeat(20_000) + "]".repeat(20_000))) { + val answer = HttpRelayHandler(MemoryRelay(backend, { VerifyPolicy }, limits = null), origins = { listOf(origin) }).ask(HttpRelayCommand.REQ, body) + assertEquals(400, answer.status, body.take(20)) + } + } + + @Test + fun aFullAuthPolicyRefusesTransportSignInUntilItOptsIn() { + val body = """{"kinds":[1]}""" + val relayUrl = RelayUrlNormalizer.normalize("wss://relay.example") + val refusing = + MemoryRelay(backend, { + object : FullAuthPolicy(relayUrl) { + override suspend fun authorize(event: RelayAuthEvent): Unit = error("backend rejected user") + } + }) + val refused = HttpRelayHandler(refusing, origins = { listOf(origin) }).ask(HttpRelayCommand.REQ, body, token(HttpRelayCommand.REQ, body)) + assertEquals(403, refused.status, refused.lines.toString()) + assertTrue(refused.lines.single().startsWith("""["CLOSED","restricted:"""), refused.lines.toString()) + + val optingIn = + MemoryRelay(backend, { + object : FullAuthPolicy(relayUrl) { + override suspend fun authorizeTransport(pubkey: HexKey): String? = null + } + }) + val signed = HttpRelayHandler(optingIn, origins = { listOf(origin) }).ask(HttpRelayCommand.REQ, body, token(HttpRelayCommand.REQ, body)) + assertEquals(200, signed.status, signed.lines.toString()) + } + + @Test + fun aFloodOfFreshTokensCannotFlushAReplay() { + val h = HttpRelayHandler(MemoryRelay(backend, false), origins = { listOf(origin) }, verifier = Nip98AuthVerifier(maxReplayEntries = 4)) + + // A token's id is its hash, so tokens signed the same second over the same body are one token: + // every one here names its own body. + fun body(n: Int) = """{"kinds":[$n]}""" + val victim = token(HttpRelayCommand.REQ, body(0)) + assertEquals(200, h.ask(HttpRelayCommand.REQ, body(0), victim).status) + for (n in 1..3) assertEquals(200, h.ask(HttpRelayCommand.REQ, body(n), token(HttpRelayCommand.REQ, body(n))).status) + val flooding = h.ask(HttpRelayCommand.REQ, body(4), token(HttpRelayCommand.REQ, body(4))) + assertEquals(429, flooding.status, "a full cache refuses the new token: ${flooding.lines}") + val replayed = h.ask(HttpRelayCommand.REQ, body(0), victim) + assertEquals(401, replayed.status, "and still remembers the old one: ${replayed.lines}") + } + + @Test + fun aMessageLimitInThePolicyChainStillRuns() { + val limited = MemoryRelay(backend, { LimitsPolicy(RelayLimits(maxMessageLength = 4096)) + VerifyPolicy }, limits = null) + val answer = HttpRelayHandler(limited, origins = { listOf(origin) }).ask(HttpRelayCommand.REQ, """{"search":"${"x".repeat(20_000)}"}""") + assertEquals(400, answer.status, answer.lines.toString()) + assertTrue(answer.lines.single().startsWith("""["NOTICE","invalid: message too large"""), answer.lines.toString()) + } + + @Test + fun aMultiByteEventUnderTheCharacterLimitIsAccepted() { + // 1,500 CJK characters: about 4,500 UTF-8 bytes, well under 4,096 characters as the engine counts. + val posted = note("\u4E2D".repeat(1_500)) + val answer = handler().ask(HttpRelayCommand.EVENT, posted.toJson()) + assertEquals(200, answer.status, answer.lines.toString()) + } + + @Test + fun aBackendFailureIsA500Line() { + val failing = + object : SessionBackend by backend { + override suspend fun sealedNegentropyStorage( + filters: List, + maxEntries: Int, + ) = error("db is down") + } + val message = Negentropy(StorageVector().also { it.seal() }, 0).initiate().toHexKey() + val answer = HttpRelayHandler(MemoryRelay(failing, { VerifyPolicy }), origins = { listOf(origin) }).ask(HttpRelayCommand.NEG, """[{"kinds":[1]},"$message"]""") + assertEquals(500, answer.status, answer.lines.toString()) + assertTrue(answer.lines.single().startsWith("""["CLOSED","error:"""), answer.lines.toString()) + } + + @Test + fun anInfiniteDeadlineStillStreams() { + backend.events += note("forever") + val answer = handler(deadline = Duration.INFINITE).ask(HttpRelayCommand.REQ, """{"kinds":[1]}""") + assertEquals(200, answer.status) + assertEquals("""["EOSE"]""", answer.lines.last()) + } + + @Test + fun aReaderStalledOnASingleAnswerIsDropped() { + val h = HttpRelayHandler(MemoryRelay(backend, false), origins = { listOf(origin) }, deadline = 100.milliseconds, tailGrace = 100.milliseconds) + val stalled = + object : HttpRelayResponse { + override suspend fun single( + status: Int, + frame: String, + ) = awaitCancellation() + + override suspend fun stream(lines: suspend HttpRelayLines.() -> Unit) = error("single") + } + assertFailsWith { + runBlocking { h.handle(HttpRelayRequest(HttpRelayCommand.COUNT, null, """{"kinds":[1]}""".encodeToByteArray()), stalled) } + } + } + + @Test + fun theBodyShapeIsReadWithoutATree() { + assertEquals("""["REQ","http",{"kinds":[1]}]""", HttpRelayCommand.REQ.frameOf(""" {"kinds":[1]} """)) + assertEquals("""["COUNT","http",{"a":"]"},{"b":"\"["}]""", HttpRelayCommand.COUNT.frameOf("""[{"a":"]"},{"b":"\"["}]""")) + assertEquals("""["NEG-OPEN","http",{"kinds":[1]},"61"]""", HttpRelayCommand.NEG.frameOf("""[{"kinds":[1]},"61"]""")) + assertEquals("""["EVENT",{"id":"x"}]""", HttpRelayCommand.EVENT.frameOf("""{"id":"x"}""")) + for (bad in listOf("", "[]", "{", "[{}", "{}}", "{]", "[{}]x", "{} {}", "[{},]", "[,{}]", "[{} {}]", "[1]", """[{},"x"]""", "null", "\"x\"")) { + assertEquals(null, HttpRelayCommand.REQ.frameOf(bad), "REQ '$bad'") + } + for (bad in listOf("""[{"kinds":[1]}]""", """["61",{}]""", """[{},61]""", """[{},"61",1]""", """{"kinds":[1]}""")) { + assertEquals(null, HttpRelayCommand.NEG.frameOf(bad), "NEG '$bad'") + } + for (bad in listOf("""[{"id":"x"}]""", "1")) assertEquals(null, HttpRelayCommand.EVENT.frameOf(bad), "EVENT '$bad'") + } + private companion object { const val STALLED_KIND = 7 const val TRICKLE_KIND = 8 From 5d67955aee14391076c0520f4bd3c76ff9797421 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 13:54:46 +0000 Subject: [PATCH 4/8] NIP-FE: remove /neg Upstream NIP-FE dropped negentropy over HTTP; NIP-77 stays on the websocket. HttpRelayCommand is REQ, COUNT and EVENT. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_016jDNhr6a3J4VC5TaG3Yd48 --- .../relay/server/backend/LiveEventStore.kt | 8 +-- .../nipFERelayOverHttp/HttpRelayCommand.kt | 16 ----- .../nipFERelayOverHttp/HttpRelayHandler.kt | 2 +- .../nipFERelayOverHttp/HttpRelayStatus.kt | 2 - .../HttpRelayHandlerTest.kt | 66 ++----------------- 5 files changed, 11 insertions(+), 83 deletions(-) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStore.kt index 1569cc194f..88da1daaed 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStore.kt @@ -409,10 +409,10 @@ class LiveEventStore( /** * Serves repeated NEG-OPENs of the same filter from one sealed storage until a write lands in its - * set. Several slots, because NIP-FE's stateless rounds re-open per round and a relay syncs more - * than one filter at a time; rebuilding costs a full scan + O(n log n) seal that grows with the - * corpus: relayBench measured 342 ms per identical-set reconcile at 50k events vs strfry's 26 ms - * off its always-current tree. + * set. Several slots, because a relay reconciles more than one filter at a time and a busy one + * takes writes between every open; rebuilding costs a full scan + O(n log n) seal that grows with + * the corpus: relayBench measured 342 ms per identical-set reconcile at 50k events vs strfry's + * 26 ms off its always-current tree. */ override suspend fun sealedNegentropyStorage( filters: List, diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayCommand.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayCommand.kt index 6d48801955..2716cba3c6 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayCommand.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayCommand.kt @@ -29,18 +29,11 @@ import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.OkMessage 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.nip77Negentropy.NegErrMessage -import com.vitorpamplona.quartz.nip77Negentropy.NegMsgMessage -import com.vitorpamplona.quartz.nip77Negentropy.NegOpenCmd /** * NIP-FE: the client commands HTTP carries, one path each. A body is the command's arguments * after its subscription id (a lone object where the command takes one); the answer ends on the * first frame [ends] accepts. - * - * `NEG` is not in NIP-FE; it is this implementation's extension: one NIP-77 round, - * `[filter, message]`. The responder keeps no state between rounds but its snapshot, which the - * backend caches per filter, so each round carries its filter and there is no session to close. */ enum class HttpRelayCommand( val path: String, @@ -48,7 +41,6 @@ enum class HttpRelayCommand( REQ("/req"), COUNT("/count"), EVENT("/event"), - NEG("/neg"), ; /** @@ -74,11 +66,6 @@ enum class HttpRelayCommand( if (!shape.isObject) return null "[\"${EventCmd.LABEL}\",${shape.text}]" } - - NEG -> { - if (shape.isObject || shape.elements != NEG_ROUND) return null - frame(NegOpenCmd.LABEL, SUB_ID, shape.inner) - } } } @@ -89,7 +76,6 @@ enum class HttpRelayCommand( REQ -> message is EoseMessage || message is ClosedMessage COUNT -> message is CountMessage || message is ClosedMessage EVENT -> message is OkMessage - NEG -> message is NegMsgMessage || message is NegErrMessage } companion object { @@ -99,8 +85,6 @@ enum class HttpRelayCommand( */ const val SUB_ID = "http" - private val NEG_ROUND = listOf('{', '"') - fun forPath(path: String): HttpRelayCommand? = entries.firstOrNull { it.path == path } private fun frame( diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt index e443028888..ef4056423d 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt @@ -346,7 +346,7 @@ class HttpRelayHandler( } /** The frames that carry a subscription id in the engine; NIP-FE sends them without it. */ -private val SUBSCRIPTION_FRAMES = setOf("EVENT", "EOSE", "CLOSED", "COUNT", "NEG-MSG", "NEG-ERR") +private val SUBSCRIPTION_FRAMES = setOf("EVENT", "EOSE", "CLOSED", "COUNT") private const val SUB_ID_FIELD = ",\"" + HttpRelayCommand.SUB_ID + "\"" diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayStatus.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayStatus.kt index 496acab463..370ca0395e 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayStatus.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayStatus.kt @@ -26,7 +26,6 @@ import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.NoticeMessage import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.OkMessage import com.vitorpamplona.quartz.nip01Core.store.RejectionReason -import com.vitorpamplona.quartz.nip77Negentropy.NegErrMessage /** * NIP-FE status codes. The status of an answer is decided by its first frame: an accepting one @@ -47,7 +46,6 @@ object HttpRelayStatus { fun of(message: Message?): Int = when (message) { is ClosedMessage -> forReason(message.message) - is NegErrMessage -> forReason(message.reason) // A duplicate is already stored, which is what the caller asked for, whichever flag the store set. is OkMessage -> if (message.success || message.message.startsWith(RejectionReason.PREFIX_DUPLICATE)) OK else forReason(message.message) is NoticeMessage -> BAD_REQUEST diff --git a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt index 4123405322..3f39895789 100644 --- a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt @@ -20,11 +20,8 @@ */ package com.vitorpamplona.quartz.nipFERelayOverHttp -import com.vitorpamplona.negentropy.Negentropy -import com.vitorpamplona.negentropy.storage.StorageVector import com.vitorpamplona.quartz.nip01Core.core.Event import com.vitorpamplona.quartz.nip01Core.core.HexKey -import com.vitorpamplona.quartz.nip01Core.core.toHexKey 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.toRelay.ReqCmd @@ -43,7 +40,6 @@ import com.vitorpamplona.quartz.nip01Core.relay.server.policies.RelayLimits import com.vitorpamplona.quartz.nip01Core.relay.server.policies.VerifyPolicy import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync import com.vitorpamplona.quartz.nip01Core.store.IEventStore -import com.vitorpamplona.quartz.nip01Core.store.IdAndTime import com.vitorpamplona.quartz.nip42RelayAuth.RelayAuthEvent import com.vitorpamplona.quartz.nip77Negentropy.NegentropySettings import com.vitorpamplona.quartz.nip98HttpAuth.HTTPAuthorizationEvent @@ -60,7 +56,7 @@ import kotlin.time.Duration import kotlin.time.Duration.Companion.milliseconds import kotlin.time.Duration.Companion.seconds -/** NIP-FE's handler over a relay engine and an in-memory backend: status, lines, auth, and negentropy rounds. */ +/** NIP-FE's handler over a relay engine and an in-memory backend: status, lines and auth. */ class HttpRelayHandlerTest { /** Everything the backend holds; a filter naming [STALLED_KIND] never answers, one naming [TRICKLE_KIND] never ends. */ private class MemoryBackend : SessionBackend { @@ -91,11 +87,6 @@ class HttpRelayHandlerTest { if (events.none { it.id == event.id }) events += event onComplete(IEventStore.InsertOutcome.Accepted) } - - override suspend fun snapshotIdsForNegentropy( - filters: List, - maxEntries: Int?, - ) = events.filter { e -> filters.any { it.match(e) } }.map { IdAndTime(it.createdAt, it.id) } } /** Refuses every read from a connection nobody signed in on, as an AUTH-gated relay does. */ @@ -238,7 +229,6 @@ class HttpRelayHandlerTest { for ((command, body) in listOf( HttpRelayCommand.REQ to "[]", HttpRelayCommand.EVENT to """[{"id":"x"}]""", - HttpRelayCommand.NEG to """[{"kinds":[1]}]""", HttpRelayCommand.COUNT to "not json", )) { val answer = handler().ask(command, body) @@ -285,45 +275,6 @@ class HttpRelayHandlerTest { assertTrue("payload" in other.lines.single(), other.lines.toString()) } - @Test - fun negentropyReconcilesInStatelessRounds() { - val shared = (1..30).map { note("shared $it", 1_700_000_000L + it) } - val onlyRelay = (1..12).map { note("relay $it", 1_700_001_000L + it) } - val onlyClient = (1..7).map { note("client $it", 1_700_002_000L + it) } - backend.events += shared + onlyRelay - - val mine = StorageVector().apply { (shared + onlyClient).forEach { insert(it.createdAt, it.id) } }.also { it.seal() } - val client = Negentropy(mine, 0) - var message = client.initiate().toHexKey() - val have = mutableSetOf() - val need = mutableSetOf() - var rounds = 0 - while (true) { - check(++rounds < 20) { "no convergence" } - // A fresh handler each round: nothing on the server side carries over. - val answer = handler().ask(HttpRelayCommand.NEG, """[{"kinds":[1]},"$message"]""") - assertEquals(200, answer.status, answer.lines.toString()) - val reply = - answer.lines - .single() - .substringAfter("""["NEG-MSG","""") - .substringBefore('"') - val result = client.reconcile(reply.hexToByteArray()) - have += result.sendIds.map { it.toHexString() } - need += result.needIds.map { it.toHexString() } - message = result.msg?.toHexKey() ?: break - } - assertEquals(onlyClient.map { it.id }.toSet(), have) - assertEquals(onlyRelay.map { it.id }.toSet(), need) - } - - @Test - fun aMalformedNegentropyRoundIsRefused() { - val answer = handler().ask(HttpRelayCommand.NEG, """[{"kinds":[1]},"zz"]""") - assertEquals(400, answer.status) - assertTrue(answer.lines.single().startsWith("""["NEG-ERR","""), answer.lines.toString()) - } - @Test fun noFirstFrameWithinTheDeadlineIsA503() { val answer = handler(deadline = 300.milliseconds).ask(HttpRelayCommand.REQ, """{"kinds":[$STALLED_KIND]}""") @@ -451,13 +402,12 @@ class HttpRelayHandlerTest { fun aBackendFailureIsA500Line() { val failing = object : SessionBackend by backend { - override suspend fun sealedNegentropyStorage( - filters: List, - maxEntries: Int, - ) = error("db is down") + override suspend fun submit( + event: Event, + onComplete: (IEventStore.InsertOutcome) -> Unit, + ): Unit = error("db is down") } - val message = Negentropy(StorageVector().also { it.seal() }, 0).initiate().toHexKey() - val answer = HttpRelayHandler(MemoryRelay(failing, { VerifyPolicy }), origins = { listOf(origin) }).ask(HttpRelayCommand.NEG, """[{"kinds":[1]},"$message"]""") + val answer = HttpRelayHandler(MemoryRelay(failing, { VerifyPolicy }), origins = { listOf(origin) }).ask(HttpRelayCommand.EVENT, note("lost").toJson()) assertEquals(500, answer.status, answer.lines.toString()) assertTrue(answer.lines.single().startsWith("""["CLOSED","error:"""), answer.lines.toString()) } @@ -491,14 +441,10 @@ class HttpRelayHandlerTest { fun theBodyShapeIsReadWithoutATree() { assertEquals("""["REQ","http",{"kinds":[1]}]""", HttpRelayCommand.REQ.frameOf(""" {"kinds":[1]} """)) assertEquals("""["COUNT","http",{"a":"]"},{"b":"\"["}]""", HttpRelayCommand.COUNT.frameOf("""[{"a":"]"},{"b":"\"["}]""")) - assertEquals("""["NEG-OPEN","http",{"kinds":[1]},"61"]""", HttpRelayCommand.NEG.frameOf("""[{"kinds":[1]},"61"]""")) assertEquals("""["EVENT",{"id":"x"}]""", HttpRelayCommand.EVENT.frameOf("""{"id":"x"}""")) for (bad in listOf("", "[]", "{", "[{}", "{}}", "{]", "[{}]x", "{} {}", "[{},]", "[,{}]", "[{} {}]", "[1]", """[{},"x"]""", "null", "\"x\"")) { assertEquals(null, HttpRelayCommand.REQ.frameOf(bad), "REQ '$bad'") } - for (bad in listOf("""[{"kinds":[1]}]""", """["61",{}]""", """[{},61]""", """[{},"61",1]""", """{"kinds":[1]}""")) { - assertEquals(null, HttpRelayCommand.NEG.frameOf(bad), "NEG '$bad'") - } for (bad in listOf("""[{"id":"x"}]""", "1")) assertEquals(null, HttpRelayCommand.EVENT.frameOf(bad), "EVENT '$bad'") } From 71d0c25b6430e4342bea6c235a05d7bf6262e4cd Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 14:00:47 +0000 Subject: [PATCH 5/8] Drop the per-filter NIP-77 snapshot cache from this PR It was here for stateless /neg rounds, which rebuilt a snapshot per round. With /neg removed it only speeds up NEG-OPEN on the websocket, which is out of scope for a NIP-FE follow-up. LiveEventStore is back to main's single-slot cache. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_016jDNhr6a3J4VC5TaG3Yd48 --- .../relay/server/backend/LiveEventStore.kt | 86 ++++------- .../LiveEventStoreSnapshotCacheTest.kt | 133 ------------------ 2 files changed, 30 insertions(+), 189 deletions(-) delete mode 100644 quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStoreSnapshotCacheTest.kt diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStore.kt index 88da1daaed..2c37a5b904 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStore.kt @@ -22,21 +22,18 @@ package com.vitorpamplona.quartz.nip01Core.relay.server.backend import com.vitorpamplona.negentropy.storage.IStorage import com.vitorpamplona.quartz.nip01Core.core.Event -import com.vitorpamplona.quartz.nip01Core.core.isAddressable -import com.vitorpamplona.quartz.nip01Core.core.isReplaceable import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.filters.FilterIndex import com.vitorpamplona.quartz.nip01Core.store.IEventStore import com.vitorpamplona.quartz.nip01Core.store.IdAndTime import com.vitorpamplona.quartz.nip01Core.store.RawEvent import com.vitorpamplona.quartz.nip01Core.store.StoreQueryContext -import com.vitorpamplona.quartz.nip09Deletions.DeletionRequestEvent -import com.vitorpamplona.quartz.nip62RequestToVanish.RequestToVanishEvent import com.vitorpamplona.quartz.utils.TimeUtils import kotlinx.coroutines.CompletableDeferred import kotlinx.coroutines.awaitCancellation import kotlinx.coroutines.withContext import kotlin.concurrent.atomics.AtomicBoolean +import kotlin.concurrent.atomics.AtomicLong import kotlin.concurrent.atomics.AtomicReference import kotlin.concurrent.atomics.ExperimentalAtomicApi @@ -157,7 +154,7 @@ class LiveEventStore( ) { ingest.submit(event, skipVerify) { outcome -> if (outcome is IEventStore.InsertOutcome.Accepted) { - forgetSnapshotsCovering(event) + writeGeneration.addAndFetch(1L) fanout(event) } onComplete(outcome) @@ -391,28 +388,32 @@ class LiveEventStore( // ------------------------------------------------------------------ /** - * The last few sealed negentropy snapshots, most recently used first, each valid until a write - * could change its set. Replaced whole on every change, never mutated, so a reader never sees a - * half-updated list. Deletion paths that bypass the ingest queue (expiration sweeps, NIP-86 admin - * purges) are not seen here, which is why entries also carry a short TTL: a snapshot is a - * point-in-time set by NIP-77's nature, and a few seconds of staleness only means a peer - * momentarily re-offers ids the relay just dropped. + * Bumped after every accepted write. A cached negentropy snapshot is + * only valid while this hasn't moved. Deletion paths that bypass the + * ingest queue (expiration sweeps, NIP-86 admin purges) don't bump it, + * which is why cache entries also carry a short TTL: a snapshot is a + * point-in-time set by NIP-77's nature, and a few seconds of staleness + * only means a peer momentarily re-offers ids the relay just dropped. */ - private val snapshotCache = AtomicReference>(emptyList()) + private val writeGeneration = AtomicLong(0L) private class CachedSnapshot( val filterKey: String, - val filters: List, + val generation: Long, val builtAt: Long, val storage: IStorage?, ) + private val snapshotCache = AtomicReference(null) + /** - * Serves repeated NEG-OPENs of the same filter from one sealed storage until a write lands in its - * set. Several slots, because a relay reconciles more than one filter at a time and a busy one - * takes writes between every open; rebuilding costs a full scan + O(n log n) seal that grows with - * the corpus: relayBench measured 342 ms per identical-set reconcile at 50k events vs strfry's - * 26 ms off its always-current tree. + * Serves repeated NEG-OPENs of the same filter from one sealed + * storage as long as no write landed in between (single slot — the + * mirror-heartbeat pattern is many peers reconciling the same broad + * filter, not many filters). Rebuilding on every open costs a full + * scan + O(n log n) seal that grows with the corpus: relayBench + * measured 342 ms per identical-set reconcile at 50k events vs + * strfry's 26 ms off its always-current tree. */ override suspend fun sealedNegentropyStorage( filters: List, @@ -427,57 +428,30 @@ class LiveEventStore( store.liveNegentropySnapshot(maxEntries)?.let { return it } } - val key = filters.joinToString(" ") { it.toJson() } + " cap=$maxEntries" + val generation = writeGeneration.load() + val key = filters.joinToString("") { it.toJson() } + "cap=$maxEntries" val now = TimeUtils.now() - snapshotCache.load().firstOrNull { it.filterKey == key && now - it.builtAt <= SNAPSHOT_TTL_SECONDS }?.let { hit -> - updateSnapshots { cached -> listOf(hit) + cached.filter { it !== hit } } - return hit.storage + val cached = snapshotCache.load() + if (cached != null && + cached.filterKey == key && + cached.generation == generation && + now - cached.builtAt <= SNAPSHOT_TTL_SECONDS + ) { + return cached.storage } val built = super.sealedNegentropyStorage(filters, maxEntries) - val entry = CachedSnapshot(key, filters, now, built) - updateSnapshots { cached -> (listOf(entry) + cached.filter { it.filterKey != key }).take(SNAPSHOT_SLOTS) } + snapshotCache.store(CachedSnapshot(key, generation, now, built)) return built } - /** - * Drops every cached snapshot whose set [event] can change: one of its filters admits the event, - * or the event can remove members it cannot see — a deletion or vanish request clears them all, - * and a replaceable or addressable event drops any snapshot whose kinds and authors it falls in, - * since the version it replaces may be in the set under tags or ids the new one no longer has. - */ - private fun forgetSnapshotsCovering(event: Event) { - if (snapshotCache.load().isEmpty()) return - if (event.kind == DeletionRequestEvent.KIND || event.kind == RequestToVanishEvent.KIND) { - updateSnapshots { emptyList() } - return - } - val supersedes = event.kind.isReplaceable() || event.kind.isAddressable() - updateSnapshots { cached -> - cached.filter { snapshot -> - snapshot.filters.none { f -> f.match(event) || (supersedes && Filter(kinds = f.kinds, authors = f.authors).match(event)) } - } - } - } - - private inline fun updateSnapshots(change: (List) -> List) { - while (true) { - val current = snapshotCache.load() - val next = change(current) - if (next == current || snapshotCache.compareAndSet(current, next)) return - } - } - private companion object { /** * Ceiling on how long a cached snapshot may serve NEG-OPENs even * with no observed writes — bounds staleness from delete paths - * the ingest queue never sees. + * the generation counter can't see. */ const val SNAPSHOT_TTL_SECONDS = 30L - - /** Distinct filters kept sealed at once; each holds up to maxSyncEvents ids at ~40 B each. */ - const val SNAPSHOT_SLOTS = 4 } } diff --git a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStoreSnapshotCacheTest.kt b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStoreSnapshotCacheTest.kt deleted file mode 100644 index eacc390d3f..0000000000 --- a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStoreSnapshotCacheTest.kt +++ /dev/null @@ -1,133 +0,0 @@ -/* - * 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.Event -import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter -import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync -import com.vitorpamplona.quartz.nip01Core.store.IEventStore -import com.vitorpamplona.quartz.nip01Core.store.IdAndTime -import com.vitorpamplona.quartz.nip01Core.store.sqlite.DefaultIndexingStrategy -import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore -import kotlinx.coroutines.CompletableDeferred -import kotlinx.coroutines.CoroutineScope -import kotlinx.coroutines.Dispatchers -import kotlinx.coroutines.SupervisorJob -import kotlinx.coroutines.cancel -import kotlinx.coroutines.runBlocking -import kotlin.test.AfterTest -import kotlin.test.Test -import kotlin.test.assertEquals - -/** The NIP-77 snapshot cache: kept across writes that cannot change its set, dropped by those that can. */ -class LiveEventStoreSnapshotCacheTest { - /** Counts how often a snapshot is actually scanned out of the store. */ - private class Counting( - private val inner: IEventStore, - ) : IEventStore by inner { - var scans = 0 - - override suspend fun snapshotIdsForNegentropy( - filters: List, - maxEntries: Int?, - onProgress: ((collected: Int) -> Unit)?, - ): List { - scans++ - return inner.snapshotIdsForNegentropy(filters, maxEntries, onProgress) - } - } - - private val scope = CoroutineScope(Dispatchers.Default + SupervisorJob()) - private val store = Counting(EventStore(dbName = null, indexStrategy = DefaultIndexingStrategy(indexEventsByPubkeyAlone = true))) - private val live = LiveEventStore(store, IngestQueue(store, scope.coroutineContext)) - private val alice = NostrSignerSync() - private var clock = 1_700_000_000L - - private val notes = listOf(Filter(kinds = listOf(1))) - private val reactions = listOf(Filter(kinds = listOf(7))) - - @AfterTest - fun tearDown() { - scope.cancel() - } - - private fun event( - kind: Int, - tags: Array> = emptyArray(), - ) = alice.sign(clock++, kind, tags, "") - - private fun publish(event: Event) = - runBlocking { - val done = CompletableDeferred() - live.submit(event) { done.complete(it) } - assertEquals(IEventStore.InsertOutcome.Accepted, done.await()) - } - - private fun open(filters: List) = runBlocking { live.sealedNegentropyStorage(filters, maxEntries = 1_000) } - - @Test - fun twoFiltersInTurnAreEachScannedOnce() { - publish(event(1)) - repeat(3) { - open(notes) - open(reactions) - } - assertEquals(2, store.scans) - } - - @Test - fun aWriteOutsideTheSetKeepsTheSnapshot() { - open(notes) - publish(event(7)) - open(notes) - assertEquals(1, store.scans) - } - - @Test - fun aWriteInsideTheSetRebuildsIt() { - open(notes) - publish(event(1)) - open(notes) - assertEquals(2, store.scans) - } - - @Test - fun aDeletionDropsEverySnapshot() { - open(notes) - open(reactions) - publish(event(5, arrayOf(arrayOf("e", "a".repeat(64))))) - open(notes) - open(reactions) - assertEquals(4, store.scans) - } - - @Test - fun aReplacementDropsASnapshotItsOldVersionWasIn() { - // The stored profile is in the set by its tag; its replacement has no such tag and so does not - // match the filter, but it removes the old version from the set all the same. - val tagged = listOf(Filter(kinds = listOf(0), tags = mapOf("t" to listOf("nostr")))) - publish(event(0, arrayOf(arrayOf("t", "nostr")))) - open(tagged) - publish(event(0)) - open(tagged) - assertEquals(2, store.scans) - } -} From 4e6472cd09decf9cfda2c5e3790e21fcc439009b Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 14:10:14 +0000 Subject: [PATCH 6/8] NIP-FE: NIP-98 tokens are not single-use Nip98AuthVerifier gains rejectReplays (default true, so NIP-86 keeps it). The NIP-FE handler turns it off: a request may land on any instance, which a per-process memory of spent tokens cannot follow, and the body's hash already limits a captured token to repeating the one command it signs inside its window. Verifier tests cover single use, a full cache, and a token without replay rejection. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_016jDNhr6a3J4VC5TaG3Yd48 --- .../quartz/nip98HttpAuth/Nip98AuthVerifier.kt | 9 +++++ .../nipFERelayOverHttp/HttpRelayHandler.kt | 15 ++++---- .../nip98HttpAuth/Nip98AuthVerifierTest.kt | 36 +++++++++++++++++++ .../HttpRelayHandlerTest.kt | 12 +++---- 4 files changed, 58 insertions(+), 14 deletions(-) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifier.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifier.kt index da83ca0d6d..7bef59da73 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifier.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifier.kt @@ -44,6 +44,7 @@ import kotlin.math.abs * 4. The `method` tag matches the HTTP method. * 5. The `u` tag matches the requested URL. * 6. If a body is present, the `payload` tag matches `sha256(body)` hex. + * 7. With [rejectReplays], the token has not been accepted before. * * Returns the verified pubkey on success; a [Result.Malformed] / * [Result.Missing] otherwise (the caller turns these into 401/403). @@ -58,6 +59,12 @@ class Nip98AuthVerifier( * captured token be replayed. Size it for the endpoint's signed-request rate over 2 x tolerance. */ private val maxReplayEntries: Int = MAX_REPLAY_ENTRIES, + /** + * Accept each token once. NIP-98 does not ask for it, and the memory is this process's alone, so + * behind a load balancer it holds per instance. Off, a token captured in its window can repeat + * the request it signs, and only that request when it binds the body's hash. + */ + private val rejectReplays: Boolean = true, ) { /** * Recently-accepted event ids → expiry epoch second. Bounded to @@ -133,6 +140,8 @@ class Nip98AuthVerifier( } } + if (!rejectReplays) return Result.Verified(event.pubKey) + // Replay check — done LAST so we don't burn a one-shot id on a // request that would otherwise have failed signature/url/etc. val expiry = nowSec + 2 * toleranceSeconds diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt index ef4056423d..fbe22d3468 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt @@ -99,8 +99,12 @@ class HttpRelayHandler( private val origins: () -> List, /** How long one answer may run, first byte to last. [Duration.INFINITE] turns the deadline off. */ private val deadline: Duration = DEFAULT_DEADLINE, - /** Its own replay cache, sized for a public endpoint, so it cannot be flushed to replay a token. */ - private val verifier: Nip98AuthVerifier = Nip98AuthVerifier(maxReplayEntries = DEFAULT_REPLAY_ENTRIES), + /** + * Tokens are not single-use: a request can land on any instance, which a per-process memory of + * spent tokens cannot follow, and the body's hash already limits a captured token to the one + * command it signs, inside its window. + */ + private val verifier: Nip98AuthVerifier = Nip98AuthVerifier(rejectReplays = false), /** Frames queued ahead of a slow reader before the answer is cut short. */ private val maxQueuedFrames: Int = DEFAULT_MAX_QUEUED_FRAMES, /** How long past the deadline the last line may take before the reader counts as stalled. */ @@ -291,7 +295,7 @@ class HttpRelayHandler( /** * A NIP-98 header, checked against the address it names when that is one of [origins], so a token - * signed at the .onion verifies there. It must bind the body's hash: it authorizes one command, once. + * signed at the .onion verifies there. It must bind the body's hash: it authorizes one command. * Another scheme (a proxy's Basic, a client's Bearer) is not addressed to the relay and is ignored. */ private suspend fun proofOf(request: HttpRelayRequest): Proof { @@ -310,7 +314,7 @@ class HttpRelayHandler( Proof.Anonymous } - // A full replay cache is the relay's limit, not the token's fault. + // A full replay cache, in a verifier that keeps one, is the relay's limit, not the token's fault. is Nip98AuthVerifier.Result.Malformed -> { if (MachineReadablePrefix.parse(r.reason) == MachineReadablePrefix.RATE_LIMITED) { Proof.Refused(r.reason) @@ -339,9 +343,6 @@ class HttpRelayHandler( const val DEFAULT_MAX_QUEUED_FRAMES = 8192 val DEFAULT_TAIL_GRACE = 5_000.milliseconds - - /** Two minutes of tokens (the replay window) at about 500 signed commands a second. */ - const val DEFAULT_REPLAY_ENTRIES = 65_536 } } diff --git a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifierTest.kt b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifierTest.kt index 2323410d5b..d92e411946 100644 --- a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifierTest.kt +++ b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifierTest.kt @@ -131,4 +131,40 @@ class Nip98AuthVerifierTest { assertTrue(r.reason.contains("kind")) } } + + @Test + fun aTokenIsAcceptedOnce() { + runBlocking { + val body = "hello".encodeToByteArray() + val (_, header) = signedToken("http://x/", "POST", body) + assertIs(verifier.verify(header, "POST", "http://x/", body)) + val again = verifier.verify(header, "POST", "http://x/", body) + assertIs(again) + assertTrue(again.reason.contains("replay")) + } + } + + @Test + fun aFullCacheRefusesTheNewTokenAndStillRemembersTheOld() { + runBlocking { + val small = Nip98AuthVerifier(now = { 1_000L }, maxReplayEntries = 2) + val tokens = (0..2).map { n -> "body $n".encodeToByteArray().let { it to signedToken("http://x/", "POST", it).second } } + for ((body, header) in tokens.take(2)) assertIs(small.verify(header, "POST", "http://x/", body)) + val (body, header) = tokens[2] + assertEquals(Nip98AuthVerifier.Result.Malformed(Nip98AuthVerifier.REPLAY_CACHE_FULL), small.verify(header, "POST", "http://x/", body)) + val (oldBody, oldHeader) = tokens[0] + assertTrue((small.verify(oldHeader, "POST", "http://x/", oldBody) as Nip98AuthVerifier.Result.Malformed).reason.contains("replay")) + } + } + + @Test + fun withoutReplayRejectionATokenVerifiesAgainButOnlyForItsBody() { + runBlocking { + val reusable = Nip98AuthVerifier(now = { 1_000L }, rejectReplays = false) + val body = "hello".encodeToByteArray() + val (_, header) = signedToken("http://x/", "POST", body) + repeat(3) { assertIs(reusable.verify(header, "POST", "http://x/", body)) } + assertIs(reusable.verify(header, "POST", "http://x/", "other".encodeToByteArray())) + } + } } diff --git a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt index 3f39895789..461012547d 100644 --- a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt @@ -262,14 +262,12 @@ class HttpRelayHandlerTest { } @Test - fun aTokenForAnotherBodyOrASecondUseIsRefused() { + fun aTokenSignsOnlyItsBodyAndIsNotSingleUse() { val h = handler() val body = """{"kinds":[1]}""" - val once = token(HttpRelayCommand.REQ, body) - assertEquals(200, h.ask(HttpRelayCommand.REQ, body, once).status) - val replayed = h.ask(HttpRelayCommand.REQ, body, once) - assertEquals(401, replayed.status) - assertTrue("replay" in replayed.lines.single(), replayed.lines.toString()) + val signed = token(HttpRelayCommand.REQ, body) + assertEquals(200, h.ask(HttpRelayCommand.REQ, body, signed).status) + assertEquals(200, h.ask(HttpRelayCommand.REQ, body, signed).status, "any instance may answer it, so none remembers it") val other = h.ask(HttpRelayCommand.REQ, """{"kinds":[0]}""", token(HttpRelayCommand.REQ, body)) assertEquals(401, other.status) assertTrue("payload" in other.lines.single(), other.lines.toString()) @@ -367,7 +365,7 @@ class HttpRelayHandlerTest { } @Test - fun aFloodOfFreshTokensCannotFlushAReplay() { + fun aHostsSingleUseVerifierThatIsFullIsA429() { val h = HttpRelayHandler(MemoryRelay(backend, false), origins = { listOf(origin) }, verifier = Nip98AuthVerifier(maxReplayEntries = 4)) // A token's id is its hash, so tokens signed the same second over the same body are one token: From 2a23450afc3d1d3c3bfc7debba54eb45edd4caba Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 14:16:44 +0000 Subject: [PATCH 7/8] NIP-FE: splice the body and let the engine parse it; verify NIP-98 against every origin JsonShape is gone. The body goes into its frame as sent and the engine parses it as it parses socket text, through the same mapper: the verb and subscription id come first and one value is read, so a body cannot become another command, and a malformed one is the engine's NOTICE. Jackson's nesting limit already bounds depth, as it does on the socket. Nip98AuthVerifier.verify takes the accepted URLs, so the handler no longer decodes the token a second time to pick one (claimedUrl). Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_016jDNhr6a3J4VC5TaG3Yd48 --- .../quartz/nip98HttpAuth/Nip98AuthVerifier.kt | 16 ++- .../nipFERelayOverHttp/HttpRelayCommand.kt | 111 ++---------------- .../nipFERelayOverHttp/HttpRelayHandler.kt | 21 +--- .../HttpRelayHandlerTest.kt | 26 +++- 4 files changed, 48 insertions(+), 126 deletions(-) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifier.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifier.kt index 7bef59da73..06d0ce31cd 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifier.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifier.kt @@ -42,7 +42,8 @@ import kotlin.math.abs * 2. Decoded body is a kind-27235 event with a valid Schnorr signature. * 3. The event's `created_at` is within ±[toleranceSeconds] of now. * 4. The `method` tag matches the HTTP method. - * 5. The `u` tag matches the requested URL. + * 5. The `u` tag is the requested URL, or one of them: a server reachable at more than one + * address (a clearnet host and an .onion) accepts a token signed at any. * 6. If a body is present, the `payload` tag matches `sha256(body)` hex. * 7. With [rejectReplays], the token has not been accepted before. * @@ -80,12 +81,19 @@ class Nip98AuthVerifier( private val seenLock = Mutex() - @OptIn(ExperimentalEncodingApi::class) suspend fun verify( authorizationHeader: String?, method: String, url: String, body: ByteArray?, + ): Result = verify(authorizationHeader, method, listOf(url), body) + + @OptIn(ExperimentalEncodingApi::class) + suspend fun verify( + authorizationHeader: String?, + method: String, + urls: Collection, + body: ByteArray?, ): Result { if (authorizationHeader.isNullOrBlank()) return Result.Missing if (!authorizationHeader.startsWith(SCHEME)) return Result.Malformed("expected '$SCHEME ' header") @@ -130,8 +138,8 @@ class Nip98AuthVerifier( if (!auth.method().equals(method, ignoreCase = true)) { return Result.Malformed("method mismatch: expected $method, got ${auth.method()}") } - if (auth.url() != url) { - return Result.Malformed("url mismatch: expected $url, got ${auth.url()}") + if (auth.url() !in urls) { + return Result.Malformed("url mismatch: expected ${urls.joinToString(" or ")}, got ${auth.url()}") } if (body != null && body.isNotEmpty()) { val expected = sha256(body).toHexKey() diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayCommand.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayCommand.kt index 2716cba3c6..cb5a865ea6 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayCommand.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayCommand.kt @@ -44,27 +44,27 @@ enum class HttpRelayCommand( ; /** - * The client frame [body] stands for, or null when it is not this command's arguments. Only - * the body's outer shape is checked here ([JsonShape]); the engine parses the frame once, so a - * malformed inside is its usual NOTICE. The shape check is what makes splicing safe: the body - * is one balanced value with nothing after it, so it cannot close the frame or open another. + * The client frame [body] stands for, or null when it plainly is not this command's arguments. + * The body is spliced in as sent and the engine parses the frame, as it parses socket text, so + * any other malformed body is the engine's NOTICE. The verb and subscription id come first and + * the parser reads one value, so nothing a body holds can make it another command. */ fun frameOf(body: String): String? { - val shape = JsonShape.of(body) ?: return null + val text = body.trim() return when (this) { REQ, COUNT -> { val filters = when { - shape.isObject -> shape.text - shape.elements.isNotEmpty() && shape.elements.all { it == '{' } -> shape.inner + text.startsWith('{') -> text + text.startsWith('[') && text.endsWith(']') -> text.substring(1, text.length - 1).trim().ifEmpty { return null } else -> return null } - frame(if (this == REQ) ReqCmd.LABEL else CountCmd.LABEL, SUB_ID, filters) + "[\"${if (this == REQ) ReqCmd.LABEL else CountCmd.LABEL}\",\"$SUB_ID\",$filters]" } EVENT -> { - if (!shape.isObject) return null - "[\"${EventCmd.LABEL}\",${shape.text}]" + if (!text.startsWith('{')) return null + "[\"${EventCmd.LABEL}\",$text]" } } } @@ -86,96 +86,5 @@ enum class HttpRelayCommand( const val SUB_ID = "http" fun forPath(path: String): HttpRelayCommand? = entries.firstOrNull { it.path == path } - - private fun frame( - label: String, - subId: String, - args: String, - ) = "[\"$label\",\"$subId\",$args]" - } -} - -/** - * The outer shape of a JSON body, read without building a tree: one object or array, brackets - * matched by type outside strings, nesting no deeper than [MAX_DEPTH], nothing after it. An array's - * [elements] are each top-level element's first character. Everything inside is left to the parser. - */ -internal class JsonShape private constructor( - val text: String, - val isObject: Boolean, - val elements: List, -) { - /** An array's contents without its brackets. */ - val inner: String get() = text.substring(1, text.length - 1) - - companion object { - /** Deep enough for any filter, event or round by a wide margin; far too shallow to exhaust a stack. */ - const val MAX_DEPTH = 32 - - fun of(body: String): JsonShape? { - val text = body.trim() - if (text.length < 2 || (text[0] != '{' && text[0] != '[')) return null - val elements = ArrayList() - val open = CharArray(MAX_DEPTH) - var depth = 0 - var inString = false - var escaped = false - // At depth 1 inside an array: whether the next non-space character starts an element. - var expectElement = text[0] == '[' - var i = 0 - while (i < text.length) { - val c = text[i] - if (inString) { - when { - escaped -> escaped = false - c == '\\' -> escaped = true - c == '"' -> inString = false - } - i++ - continue - } - if (depth == 0 && i > 0) return null - if (depth == 1 && text[0] == '[' && !c.isWhitespace()) { - when { - c == ',' -> { - if (expectElement) return null - expectElement = true - i++ - continue - } - - c == ']' -> { - if (expectElement && elements.isNotEmpty()) return null - } - - expectElement -> { - elements.add(c) - expectElement = false - } - - c == '{' || c == '[' || c == '"' -> { - return null - } - } - } - when (c) { - '"' -> { - inString = true - } - - '{', '[' -> { - if (depth == MAX_DEPTH) return null - open[depth++] = c - } - - '}', ']' -> { - if (depth == 0 || open[--depth] != (if (c == '}') '{' else '[')) return null - } - } - i++ - } - if (depth != 0 || inString) return null - return JsonShape(text, text[0] == '{', elements) - } } } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt index fbe22d3468..90b611720e 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt @@ -21,14 +21,12 @@ package com.vitorpamplona.quartz.nipFERelayOverHttp import com.vitorpamplona.quartz.nip01Core.core.HexKey -import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper 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.MachineReadablePrefix import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message import com.vitorpamplona.quartz.nip01Core.relay.server.RelayServerBase import com.vitorpamplona.quartz.nip01Core.relay.server.SessionSink -import com.vitorpamplona.quartz.nip98HttpAuth.HTTPAuthorizationEvent import com.vitorpamplona.quartz.nip98HttpAuth.Nip98AuthVerifier import kotlinx.coroutines.CancellationException import kotlinx.coroutines.CompletableDeferred @@ -38,8 +36,6 @@ import kotlinx.coroutines.coroutineScope import kotlinx.coroutines.launch import kotlinx.coroutines.withTimeout import kotlinx.coroutines.withTimeoutOrNull -import kotlin.io.encoding.Base64 -import kotlin.io.encoding.ExperimentalEncodingApi import kotlin.time.Duration import kotlin.time.Duration.Companion.milliseconds import kotlin.time.TimeSource @@ -294,8 +290,8 @@ class HttpRelayHandler( } /** - * A NIP-98 header, checked against the address it names when that is one of [origins], so a token - * signed at the .onion verifies there. It must bind the body's hash: it authorizes one command. + * A NIP-98 header, checked against every address in [origins], so a token signed at the .onion + * verifies there. It must bind the body's hash: it authorizes one command. * Another scheme (a proxy's Basic, a client's Bearer) is not addressed to the relay and is ignored. */ private suspend fun proofOf(request: HttpRelayRequest): Proof { @@ -304,8 +300,8 @@ class HttpRelayHandler( if (!header.regionMatches(0, scheme, 0, scheme.length, ignoreCase = true)) return Proof.Anonymous val token = scheme + header.substring(scheme.length).trim() val accepted = origins().map { it.trimEnd('/') + request.command.path } - val url = claimedUrl(token)?.takeIf { it in accepted } ?: accepted.firstOrNull() ?: return Proof.Refused(MachineReadablePrefix.AUTH_REQUIRED.format("this relay names no url to sign")) - return when (val r = verifier.verify(token, "POST", url, request.body)) { + if (accepted.isEmpty()) return Proof.Refused(MachineReadablePrefix.AUTH_REQUIRED.format("this relay names no url to sign")) + return when (val r = verifier.verify(token, "POST", accepted, request.body)) { is Nip98AuthVerifier.Result.Verified -> { Proof.Signed(r.pubkey) } @@ -325,15 +321,6 @@ class HttpRelayHandler( } } - /** The `u` a NIP-98 token names, read as the verifier reads it, or null when it does not decode. */ - @OptIn(ExperimentalEncodingApi::class) - private fun claimedUrl(token: String): String? = - runCatching { - val json = Base64.decode(token.removePrefix(Nip98AuthVerifier.SCHEME).trim()).decodeToString() - val event = OptimizedJsonMapper.fromJson(json) - HTTPAuthorizationEvent(event.id, event.pubKey, event.createdAt, event.tags, event.content, event.sig).url() - }.getOrNull() - private fun closed(reason: String) = withoutSubId(ClosedMessage(HttpRelayCommand.SUB_ID, reason).toJson()) companion object { diff --git a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt index 461012547d..1c42b052eb 100644 --- a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt @@ -436,16 +436,34 @@ class HttpRelayHandlerTest { } @Test - fun theBodyShapeIsReadWithoutATree() { + fun aBodyIsSplicedIntoItsFrameAsSent() { assertEquals("""["REQ","http",{"kinds":[1]}]""", HttpRelayCommand.REQ.frameOf(""" {"kinds":[1]} """)) assertEquals("""["COUNT","http",{"a":"]"},{"b":"\"["}]""", HttpRelayCommand.COUNT.frameOf("""[{"a":"]"},{"b":"\"["}]""")) assertEquals("""["EVENT",{"id":"x"}]""", HttpRelayCommand.EVENT.frameOf("""{"id":"x"}""")) - for (bad in listOf("", "[]", "{", "[{}", "{}}", "{]", "[{}]x", "{} {}", "[{},]", "[,{}]", "[{} {}]", "[1]", """[{},"x"]""", "null", "\"x\"")) { - assertEquals(null, HttpRelayCommand.REQ.frameOf(bad), "REQ '$bad'") - } + for (bad in listOf("", "[]", "[ ]", "[{}", "1", "null", "\"x\"")) assertEquals(null, HttpRelayCommand.REQ.frameOf(bad), "REQ '$bad'") for (bad in listOf("""[{"id":"x"}]""", "1")) assertEquals(null, HttpRelayCommand.EVENT.frameOf(bad), "EVENT '$bad'") } + @Test + fun aBodyCannotCarryASecondCommand() { + val smuggled = note("smuggled") + val answer = handler().ask(HttpRelayCommand.REQ, """{"kinds":[1]}],["EVENT",${smuggled.toJson()}""") + assertTrue(answer.lines.last().let { it == """["EOSE"]""" || it.startsWith("""["NOTICE",""") }, answer.lines.toString()) + assertTrue(backend.events.none { it.id == smuggled.id }, "only the REQ ran") + } + + @Test + fun aTokenSignedAtAnyOfTheRelaysAddressesVerifies() { + val onion = "http://relayxyz.onion" + val h = HttpRelayHandler(MemoryRelay(backend, true), origins = { listOf(origin, onion) }) + val body = """{"kinds":[1]}""" + + fun at(base: String) = alice.sign(HTTPAuthorizationEvent.build(base + "/req", "POST", body.encodeToByteArray(), System.currentTimeMillis() / 1000) {}).toAuthToken() + assertEquals(200, h.ask(HttpRelayCommand.REQ, body, at(onion)).status) + assertEquals(200, h.ask(HttpRelayCommand.REQ, body, at(origin)).status) + assertEquals(401, h.ask(HttpRelayCommand.REQ, body, at("https://elsewhere.example")).status) + } + private companion object { const val STALLED_KIND = 7 const val TRICKLE_KIND = 8 From 2442a12338acb5f7159bd38a89d5ab0c74cbf965 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 14:32:33 +0000 Subject: [PATCH 8/8] Leave Nip98AuthVerifier as on main NIP-86 keeps the verifier exactly as it was: no rejectReplays, no maxReplayEntries, no URL-list overload. The NIP-FE handler gets non-single-use tokens from a fresh verifier per request, trying each of the relay's addresses in turn, and loses its verifier parameter. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_016jDNhr6a3J4VC5TaG3Yd48 --- .../quartz/nip98HttpAuth/Nip98AuthVerifier.kt | 56 ++++++------------- .../nipFERelayOverHttp/HttpRelayHandler.kt | 32 ++++------- .../nip98HttpAuth/Nip98AuthVerifierTest.kt | 36 ------------ .../HttpRelayHandlerTest.kt | 17 ------ 4 files changed, 28 insertions(+), 113 deletions(-) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifier.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifier.kt index 06d0ce31cd..ebd580ee15 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifier.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifier.kt @@ -42,10 +42,8 @@ import kotlin.math.abs * 2. Decoded body is a kind-27235 event with a valid Schnorr signature. * 3. The event's `created_at` is within ±[toleranceSeconds] of now. * 4. The `method` tag matches the HTTP method. - * 5. The `u` tag is the requested URL, or one of them: a server reachable at more than one - * address (a clearnet host and an .onion) accepts a token signed at any. + * 5. The `u` tag matches the requested URL. * 6. If a body is present, the `payload` tag matches `sha256(body)` hex. - * 7. With [rejectReplays], the token has not been accepted before. * * Returns the verified pubkey on success; a [Result.Malformed] / * [Result.Missing] otherwise (the caller turns these into 401/403). @@ -54,22 +52,10 @@ class Nip98AuthVerifier( private val now: () -> Long = { TimeUtils.now() }, /** Allowed clock skew in seconds. NIP-98 says 60. */ private val toleranceSeconds: Long = 60, - /** - * Tokens remembered at once. When every remembered token is still inside its window, a new one - * is refused (`rate-limited:`) rather than an unexpired one forgotten: forgetting is what lets a - * captured token be replayed. Size it for the endpoint's signed-request rate over 2 x tolerance. - */ - private val maxReplayEntries: Int = MAX_REPLAY_ENTRIES, - /** - * Accept each token once. NIP-98 does not ask for it, and the memory is this process's alone, so - * behind a load balancer it holds per instance. Off, a token captured in its window can repeat - * the request it signs, and only that request when it binds the body's hash. - */ - private val rejectReplays: Boolean = true, ) { /** * Recently-accepted event ids → expiry epoch second. Bounded to - * [maxReplayEntries], never by evicting a live entry; each entry expires + * [MAX_REPLAY_ENTRIES] (insertion-order eviction); each entry expires * after `2 × toleranceSeconds` (twice the accepted window so a token * can't be reused by an attacker who buffers across the boundary). * @@ -81,18 +67,11 @@ class Nip98AuthVerifier( private val seenLock = Mutex() - suspend fun verify( - authorizationHeader: String?, - method: String, - url: String, - body: ByteArray?, - ): Result = verify(authorizationHeader, method, listOf(url), body) - @OptIn(ExperimentalEncodingApi::class) suspend fun verify( authorizationHeader: String?, method: String, - urls: Collection, + url: String, body: ByteArray?, ): Result { if (authorizationHeader.isNullOrBlank()) return Result.Missing @@ -138,8 +117,8 @@ class Nip98AuthVerifier( if (!auth.method().equals(method, ignoreCase = true)) { return Result.Malformed("method mismatch: expected $method, got ${auth.method()}") } - if (auth.url() !in urls) { - return Result.Malformed("url mismatch: expected ${urls.joinToString(" or ")}, got ${auth.url()}") + if (auth.url() != url) { + return Result.Malformed("url mismatch: expected $url, got ${auth.url()}") } if (body != null && body.isNotEmpty()) { val expected = sha256(body).toHexKey() @@ -148,8 +127,6 @@ class Nip98AuthVerifier( } } - if (!rejectReplays) return Result.Verified(event.pubKey) - // Replay check — done LAST so we don't burn a one-shot id on a // request that would otherwise have failed signature/url/etc. val expiry = nowSec + 2 * toleranceSeconds @@ -163,13 +140,18 @@ class Nip98AuthVerifier( while (it.hasNext()) { if (it.next().value <= nowSec) it.remove() else break } - if (seenEventIds.containsKey(event.id)) { + if (seenEventIds.put(event.id, expiry) != null) { return Result.Malformed("replay: this NIP-98 token has already been used") } - if (seenEventIds.size >= maxReplayEntries) { - return Result.Malformed(REPLAY_CACHE_FULL) + // Cap entries: drop oldest by insertion order. Equivalent to + // the JDK LinkedHashMap.removeEldestEntry hook we used before, + // but works in KMP commonMain. + while (seenEventIds.size > MAX_REPLAY_ENTRIES) { + val eldest = seenEventIds.keys.iterator() + if (!eldest.hasNext()) break + eldest.next() + eldest.remove() } - seenEventIds[event.id] = expiry } return Result.Verified(event.pubKey) @@ -191,13 +173,11 @@ class Nip98AuthVerifier( const val SCHEME = "Nostr " /** - * Default cap on the in-memory replay cache: with a 60s tolerance, 1024 live tokens is about - * 8 signed requests a second, generous for an admin endpoint. A public endpoint passes a - * larger `maxReplayEntries`. + * Cap on the in-memory replay-cache size. With a 60s tolerance + * an attacker would need to push >MAX/120 verified requests per + * second (one new id per ~120 ms) to evict legitimate entries. + * 1024 is generous for an admin endpoint. */ const val MAX_REPLAY_ENTRIES = 1024 - - /** Answered when the replay cache holds nothing but live tokens. */ - const val REPLAY_CACHE_FULL = "rate-limited: too many fresh NIP-98 tokens at once; retry shortly" } } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt index 90b611720e..a4054bbd08 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt @@ -95,12 +95,6 @@ class HttpRelayHandler( private val origins: () -> List, /** How long one answer may run, first byte to last. [Duration.INFINITE] turns the deadline off. */ private val deadline: Duration = DEFAULT_DEADLINE, - /** - * Tokens are not single-use: a request can land on any instance, which a per-process memory of - * spent tokens cannot follow, and the body's hash already limits a captured token to the one - * command it signs, inside its window. - */ - private val verifier: Nip98AuthVerifier = Nip98AuthVerifier(rejectReplays = false), /** Frames queued ahead of a slow reader before the answer is cut short. */ private val maxQueuedFrames: Int = DEFAULT_MAX_QUEUED_FRAMES, /** How long past the deadline the last line may take before the reader counts as stalled. */ @@ -301,24 +295,18 @@ class HttpRelayHandler( val token = scheme + header.substring(scheme.length).trim() val accepted = origins().map { it.trimEnd('/') + request.command.path } if (accepted.isEmpty()) return Proof.Refused(MachineReadablePrefix.AUTH_REQUIRED.format("this relay names no url to sign")) - return when (val r = verifier.verify(token, "POST", accepted, request.body)) { - is Nip98AuthVerifier.Result.Verified -> { - Proof.Signed(r.pubkey) - } - - is Nip98AuthVerifier.Result.Missing -> { - Proof.Anonymous - } - - // A full replay cache, in a verifier that keeps one, is the relay's limit, not the token's fault. - is Nip98AuthVerifier.Result.Malformed -> { - if (MachineReadablePrefix.parse(r.reason) == MachineReadablePrefix.RATE_LIMITED) { - Proof.Refused(r.reason) - } else { - Proof.Refused(MachineReadablePrefix.AUTH_REQUIRED.format("NIP-98 ${r.reason}")) - } + // A fresh verifier each time: it remembers the tokens it accepts, and a NIP-FE token is not + // single-use. A request can land on any instance, which no one process's memory can follow, + // and the body's hash already limits a captured token to the command it signs, in its window. + var refusal: Nip98AuthVerifier.Result.Malformed? = null + for (url in accepted) { + when (val r = Nip98AuthVerifier().verify(token, "POST", url, request.body)) { + is Nip98AuthVerifier.Result.Verified -> return Proof.Signed(r.pubkey) + is Nip98AuthVerifier.Result.Missing -> return Proof.Anonymous + is Nip98AuthVerifier.Result.Malformed -> refusal = refusal ?: r } } + return Proof.Refused(MachineReadablePrefix.AUTH_REQUIRED.format("NIP-98 ${refusal?.reason}")) } private fun closed(reason: String) = withoutSubId(ClosedMessage(HttpRelayCommand.SUB_ID, reason).toJson()) diff --git a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifierTest.kt b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifierTest.kt index d92e411946..2323410d5b 100644 --- a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifierTest.kt +++ b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/nip98HttpAuth/Nip98AuthVerifierTest.kt @@ -131,40 +131,4 @@ class Nip98AuthVerifierTest { assertTrue(r.reason.contains("kind")) } } - - @Test - fun aTokenIsAcceptedOnce() { - runBlocking { - val body = "hello".encodeToByteArray() - val (_, header) = signedToken("http://x/", "POST", body) - assertIs(verifier.verify(header, "POST", "http://x/", body)) - val again = verifier.verify(header, "POST", "http://x/", body) - assertIs(again) - assertTrue(again.reason.contains("replay")) - } - } - - @Test - fun aFullCacheRefusesTheNewTokenAndStillRemembersTheOld() { - runBlocking { - val small = Nip98AuthVerifier(now = { 1_000L }, maxReplayEntries = 2) - val tokens = (0..2).map { n -> "body $n".encodeToByteArray().let { it to signedToken("http://x/", "POST", it).second } } - for ((body, header) in tokens.take(2)) assertIs(small.verify(header, "POST", "http://x/", body)) - val (body, header) = tokens[2] - assertEquals(Nip98AuthVerifier.Result.Malformed(Nip98AuthVerifier.REPLAY_CACHE_FULL), small.verify(header, "POST", "http://x/", body)) - val (oldBody, oldHeader) = tokens[0] - assertTrue((small.verify(oldHeader, "POST", "http://x/", oldBody) as Nip98AuthVerifier.Result.Malformed).reason.contains("replay")) - } - } - - @Test - fun withoutReplayRejectionATokenVerifiesAgainButOnlyForItsBody() { - runBlocking { - val reusable = Nip98AuthVerifier(now = { 1_000L }, rejectReplays = false) - val body = "hello".encodeToByteArray() - val (_, header) = signedToken("http://x/", "POST", body) - repeat(3) { assertIs(reusable.verify(header, "POST", "http://x/", body)) } - assertIs(reusable.verify(header, "POST", "http://x/", "other".encodeToByteArray())) - } - } } diff --git a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt index 1c42b052eb..31695d8da3 100644 --- a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt @@ -43,7 +43,6 @@ import com.vitorpamplona.quartz.nip01Core.store.IEventStore import com.vitorpamplona.quartz.nip42RelayAuth.RelayAuthEvent import com.vitorpamplona.quartz.nip77Negentropy.NegentropySettings import com.vitorpamplona.quartz.nip98HttpAuth.HTTPAuthorizationEvent -import com.vitorpamplona.quartz.nip98HttpAuth.Nip98AuthVerifier import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.awaitCancellation import kotlinx.coroutines.runBlocking @@ -364,22 +363,6 @@ class HttpRelayHandlerTest { assertEquals(200, signed.status, signed.lines.toString()) } - @Test - fun aHostsSingleUseVerifierThatIsFullIsA429() { - val h = HttpRelayHandler(MemoryRelay(backend, false), origins = { listOf(origin) }, verifier = Nip98AuthVerifier(maxReplayEntries = 4)) - - // A token's id is its hash, so tokens signed the same second over the same body are one token: - // every one here names its own body. - fun body(n: Int) = """{"kinds":[$n]}""" - val victim = token(HttpRelayCommand.REQ, body(0)) - assertEquals(200, h.ask(HttpRelayCommand.REQ, body(0), victim).status) - for (n in 1..3) assertEquals(200, h.ask(HttpRelayCommand.REQ, body(n), token(HttpRelayCommand.REQ, body(n))).status) - val flooding = h.ask(HttpRelayCommand.REQ, body(4), token(HttpRelayCommand.REQ, body(4))) - assertEquals(429, flooding.status, "a full cache refuses the new token: ${flooding.lines}") - val replayed = h.ask(HttpRelayCommand.REQ, body(0), victim) - assertEquals(401, replayed.status, "and still remembers the old one: ${replayed.lines}") - } - @Test fun aMessageLimitInThePolicyChainStillRuns() { val limited = MemoryRelay(backend, { LimitsPolicy(RelayLimits(maxMessageLength = 4096)) + VerifyPolicy }, limits = null)