feat(quartz): NIP-FE relay commands over HTTP, and the engine seams it needs

A relay can now answer one client command per HTTP request (POST /req,
/count, /event, /neg) with its own NIP-01 frames, one per line, ending
on the command's answer, through the same session, policies and backend
as its websocket. The spec draft is NIP-FE, "Relay Commands over HTTP".

Engine (nip01Core/relay/server), source- and binary-compatible:

- SessionSink: where a session's frames go, typed (Message) wherever the
  engine built one and raw for the spliced EVENT frames. A transport can
  now end an answer at its EOSE or map a CLOSED onto a status without
  re-parsing wire JSON. SessionSink.of(send) is the old string callback,
  and RelaySession/RelayServerBase keep every (String) -> Unit overload.
- RelaySession.initialAuthenticatedUsers, RelayServerBase.connect(sink,
  authenticatedUsers) and serve(sink, authenticatedUsers, incoming): a
  transport that proved a key itself (NIP-98) opens the session signed
  in, recorded exactly where a NIP-42 AUTH would record it, so every
  policy and backend that reads RequestContext.authenticatedUsers sees
  it with no side table.
- RelaySession.receive(Command): dispatch an already-parsed command. The
  string overload parses and calls it; acceptMessage still guards text.

NIP-FE (nipFERelayOverHttp), common code, no HTTP library:

- HttpRelayCommand: the four paths, each body parsed into its typed
  Command through a JSON tree (so a body can only ever be arguments),
  and which Message ends each answer.
- HttpRelayStatus: the status an answer's first frame gives it, from
  its NIP-01 prefix; a duplicate: OK is a 200.
- HttpRelayHandler: size, parse, NIP-98 (payload-bound, single-use, its
  own replay cache, any of the relay's addresses, other schemes
  ignored), then the exchange: first frame decides the status, the rest
  streams through HttpRelayResponse/HttpRelayLines, flushed whenever no
  frame is waiting. The deadline is read between frames and never
  interrupts a write; a reader that stops reading is HttpRelayReaderStalled
  past a grace, and the host drops it. Admission stays with the host, in
  front of handle(), so a refused request never spends a token.
- NEG is one stateless NIP-77 round, [filter, message]: the responder's
  only state is its sealed snapshot, which the backend already caches
  per filter, so nothing is held between rounds and no NEG-CLOSE exists.

Not in this change: backpressure from the socket into the backend. The
store callbacks cannot suspend, so a reader more than maxQueuedFrames
behind is cut with a CLOSED line, the websocket's bound.

Verified with :quartz:jvmTest (HttpRelayHandlerTest, 13 cases over a fake
in-memory SessionBackend), and earlier against the a8e8778265 jar that
vespa-relay pins, where every call the rest of the jar makes into
RelaySession/RelayServerBase still links.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_016jDNhr6a3J4VC5TaG3Yd48
This commit is contained in:
Claude
2026-09-27 01:38:37 +00:00
parent b1b3aca5c4
commit 400c15ff44
7 changed files with 974 additions and 10 deletions
@@ -20,6 +20,7 @@
*/
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
@@ -79,15 +80,27 @@ abstract class RelayServerBase(
* @param send Callback the server uses to send JSON messages to this client.
* Implementations must be safe to call from any coroutine.
*/
fun connect(send: (String) -> Unit): RelaySession =
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<HexKey> = emptySet(),
): RelaySession =
connections.register(
RelaySession(
policy = buildPolicy(),
store = backend,
scope = scope,
onSend = send,
sink = sink,
onClose = { connections.unregister(it.id) },
negentropySettings = negentropySettings,
initialAuthenticatedUsers = authenticatedUsers,
),
)
@@ -102,8 +115,15 @@ abstract class RelayServerBase(
suspend fun serve(
send: (String) -> Unit,
incoming: suspend (RelaySession) -> Unit,
) = serve(SessionSink.of(send), emptySet(), incoming)
/** [serve] over a [SessionSink], signed in as [authenticatedUsers] from the start. */
suspend fun serve(
sink: SessionSink,
authenticatedUsers: Set<HexKey>,
incoming: suspend (RelaySession) -> Unit,
) {
val session = connect(send)
val session = connect(sink, authenticatedUsers)
try {
incoming(session)
} finally {
@@ -32,6 +32,7 @@ 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.AuthCmd
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.CloseCmd
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
@@ -67,7 +68,8 @@ class RelaySession(
private val store: SessionBackend,
val policy: IRelayPolicy,
private val scope: CoroutineScope,
private val onSend: (String) -> Unit,
/** Where every frame for this client goes; see [SessionSink] for the typed/raw split. */
private val sink: SessionSink,
private val onClose: (RelaySession) -> Unit,
negentropySettings: NegentropySettings = NegentropySettings.Default,
/**
@@ -77,7 +79,26 @@ 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<HexKey> = emptySet(),
) : AutoCloseable {
/** The original, string-only constructor; every frame goes to [onSend] as wire JSON. */
constructor(
store: SessionBackend,
policy: IRelayPolicy,
scope: CoroutineScope,
onSend: (String) -> Unit,
onClose: (RelaySession) -> Unit,
negentropySettings: NegentropySettings = NegentropySettings.Default,
id: Long = nextConnectionId(),
) : this(store, policy, scope, SessionSink.of(onSend), onClose, negentropySettings, id)
private val subscriptions = LargeCache<String, Job>()
/**
@@ -95,7 +116,7 @@ class RelaySession(
* the copy costs nothing on the hot path.
*/
@Volatile
private var authenticatedUsers = setOf<HexKey>()
private var authenticatedUsers: Set<HexKey> = initialAuthenticatedUsers.toSet()
/**
* The per-connection scope. Handed to the [policy] at connect (so gating
@@ -133,10 +154,7 @@ class RelaySession(
fun send(message: Message) {
try {
// message.toJson() defaults to OptimizedJsonMapper.toJson(this) for
// every type; NegMsgMessage overrides it with a direct-build wire
// path (identical output, ~2× faster on big reconcile frames).
onSend(message.toJson())
sink.message(message)
} catch (e: Exception) {
Log.w("ClientSession") { "Failed to send to ${e.message}" }
}
@@ -145,7 +163,7 @@ class RelaySession(
/** [send] for frames that are already wire-format JSON (the raw REQ path). */
private fun sendRaw(json: String) {
try {
onSend(json)
sink.raw(json)
} catch (e: Exception) {
Log.w("ClientSession") { "Failed to send to ${e.message}" }
}
@@ -175,6 +193,16 @@ class RelaySession(
return
}
receive(cmd)
}
/**
* Dispatches an already-parsed command, for a transport that parsed and
* sized it itself (NIP-FE's HTTP bodies). [IRelayPolicy.acceptMessage]
* is not run here — it judges wire text — so such a caller applies the
* message-length limit before calling.
*/
suspend fun receive(cmd: Command) {
if (!cmd.isValid()) {
send(NoticeMessage("error: invalid command"))
return
@@ -0,0 +1,54 @@
/*
* Copyright (c) 2025 Vitor Pamplona
*
* Permission is hereby granted, free of charge, to any person obtaining a copy of
* this software and associated documentation files (the "Software"), to deal in
* the Software without restriction, including without limitation the rights to use,
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
* Software, and to permit persons to whom the Software is furnished to do so,
* subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all
* copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
*/
package com.vitorpamplona.quartz.nip01Core.relay.server
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message
/**
* Where a [RelaySession]'s frames go. The engine hands over a [Message]
* wherever it has one, so a transport that must react to a frame (end an
* HTTP answer at its EOSE, map a CLOSED onto a status) reads the type
* instead of re-parsing wire JSON; stored and live EVENT frames, which the
* zero-decode REQ path builds as text, arrive through [raw].
*
* Both are called from the engine's coroutines, on any thread, and must not
* block. Deliberately not a `fun interface`: a lambda would be ambiguous with
* the `(String) -> Unit` overloads that predate it.
*/
interface SessionSink {
/** A frame the engine built as a [Message]. */
fun message(message: Message)
/** A frame already in wire JSON: an `EVENT` spliced from stored or live text. */
fun raw(json: String)
companion object {
/** Every frame to [send] as wire JSON, which is all a socket wants. */
fun of(send: (String) -> Unit): SessionSink =
object : SessionSink {
// NegMsgMessage overrides toJson with a direct-build wire path
// (identical output, ~2× faster on big reconcile frames).
override fun message(message: Message) = send(message.toJson())
override fun raw(json: String) = send(json)
}
}
}
@@ -0,0 +1,139 @@
/*
* 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.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.
*
* `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.
*/
enum class HttpRelayCommand(
val path: String,
) {
REQ("/req"),
COUNT("/count"),
EVENT("/event"),
NEG("/neg"),
;
/**
* 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.
*/
fun parse(body: String): Parsed? {
val tree =
try {
Json.parseToJsonElement(body)
} catch (_: SerializationException) {
return null
}
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,
)
/** Whether [message] is the last frame of this command's answer. */
fun ends(message: Message): Boolean =
message is NoticeMessage ||
when (this) {
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
}
private fun head(): List<JsonElement> =
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<JsonElement>? =
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. */
const val SUB_ID = "http"
fun forPath(path: String): HttpRelayCommand? = entries.firstOrNull { it.path == path }
}
}
@@ -0,0 +1,295 @@
/*
* 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.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.Nip98AuthVerifier
import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.TimeoutCancellationException
import kotlinx.coroutines.channels.Channel
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.TimeSource
/** One NIP-FE request as the handler needs it. The host routes by [HttpRelayCommand.path] and bounds the body while reading it. */
class HttpRelayRequest(
val command: HttpRelayCommand,
/** The `Authorization` header as sent, or null. */
val authorization: String?,
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`.
*/
interface HttpRelayResponse {
/** An answer that is one frame, with its status: a refusal, or a command answered at once. */
suspend fun single(
status: Int,
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.
*/
suspend fun stream(lines: suspend HttpRelayLines.() -> Unit)
}
/** The body of a streamed answer. */
interface HttpRelayLines {
/** One frame and its newline. May suspend on the socket: that is the backpressure. */
suspend fun line(frame: String)
/** Sends what [line] has written so far; called whenever no frame is waiting. */
suspend fun flush()
}
/** The client stopped reading a streamed 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.
*/
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<String>,
/** 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(),
/** 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,
) {
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"))
}
val parsed =
command.parse(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) {
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)))
}
exchange(parsed, command, signedIn, response)
}
/** A frame as queued: its wire text, and its type when the engine built one. */
private class Frame(
val json: String,
val message: Message?,
)
private fun HttpRelayCommand.ends(frame: Frame) = frame.message?.let(::ends) == true
private suspend fun exchange(
parsed: HttpRelayCommand.Parsed,
command: HttpRelayCommand,
signedIn: Set<HexKey>,
response: HttpRelayResponse,
) = coroutineScope {
val frames = Channel<Frame>(maxQueuedFrames)
val ended = CompletableDeferred<Unit>()
// 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,
) {
val sent = frames.trySend(frame)
if (sent.isClosed) return
if (sent.isFailure || last) {
frames.close()
ended.complete(Unit)
}
}
val sink =
object : SessionSink {
override fun message(message: Message) {
// The challenge every connection opens with; this one proved its key by NIP-98 instead.
if (message is AuthMessage) return
offer(Frame(message.toJson(), message), command.ends(message))
}
override fun raw(json: String) = offer(Frame(json, null), false)
}
val session =
launch {
server.serve(sink, signedIn) {
it.receive(parsed.command)
ended.await()
}
}
val started = TimeSource.Monotonic.markNow()
fun remainingMs() = (deadlineMs - started.elapsedNow().inWholeMilliseconds).coerceAtLeast(0)
try {
val first = withTimeoutOrNull(deadlineMs) { frames.receiveCatching().getOrNull() }
when {
first == null -> {
response.single(HttpRelayStatus.UNAVAILABLE, closed("error: no answer within ${deadlineMs / 1000}s"))
}
command.ends(first) || HttpRelayStatus.of(first.message) != HttpRelayStatus.OK -> {
response.single(HttpRelayStatus.of(first.message), 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()
}
} catch (_: TimeoutCancellationException) {
throw HttpRelayReaderStalled()
}
}
}
}
} finally {
session.cancel()
}
}
/** 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.
*/
private suspend fun HttpRelayLines.drain(
first: Frame,
frames: Channel<Frame>,
command: HttpRelayCommand,
started: TimeSource.Monotonic.ValueTimeMark,
): Ending {
var frame = first
while (true) {
line(frame.json)
if (command.ends(frame)) 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
// 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
}
}
/** Who a request acts as. */
private sealed interface Proof {
data object Anonymous : Proof
class Signed(
val pubkey: HexKey,
) : Proof
class Refused(
val reason: String,
) : Proof
}
/**
* 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.
* 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 {
val header = request.authorization?.trim().orEmpty()
val scheme = Nip98AuthVerifier.SCHEME
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")
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
}
}
/** The `u` tag of a NIP-98 token, or null when it does not decode; the verifier then says why. */
@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)
}.getOrNull()
private fun closed(reason: String) = ClosedMessage(HttpRelayCommand.SUB_ID, reason).toJson()
companion object {
const val DEFAULT_DEADLINE_MS = 30_000L
/** 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
}
}
@@ -0,0 +1,65 @@
/*
* 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.nipFERelayOverHttp
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.ClosedMessage
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.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.
*/
object HttpRelayStatus {
const val OK = 200
const val BAD_REQUEST = 400
const val UNAUTHORIZED = 401
const val FORBIDDEN = 403
const val PAYLOAD_TOO_LARGE = 413
const val TOO_MANY_REQUESTS = 429
const val INTERNAL_ERROR = 500
const val UNAVAILABLE = 503
/** The status an answer opening with [message] carries. A raw EVENT frame 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 NoticeMessage -> BAD_REQUEST
else -> OK
}
/** The status a NIP-01 machine-readable prefix stands for. */
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
}
}
@@ -0,0 +1,363 @@
/*
* 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.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.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.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.PassThroughPolicy
import com.vitorpamplona.quartz.nip01Core.relay.server.policies.PolicyResult
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.nip77Negentropy.NegentropySettings
import com.vitorpamplona.quartz.nip98HttpAuth.HTTPAuthorizationEvent
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.awaitCancellation
import kotlinx.coroutines.runBlocking
import java.util.concurrent.CopyOnWriteArrayList
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFailsWith
import kotlin.test.assertTrue
/** NIP-FE's handler over a relay engine and an in-memory backend: status, lines, auth, and negentropy rounds. */
class HttpRelayHandlerTest {
/** Everything the backend holds; a filter naming [STALLED_KIND] never answers, one naming [TRICKLE_KIND] never ends. */
private class MemoryBackend : SessionBackend {
val events = CopyOnWriteArrayList<Event>()
override suspend fun query(
ctx: RequestContext,
filters: List<Filter>,
onEach: (Event) -> Unit,
onEose: () -> Unit,
) {
if (filters.any { it.kinds?.contains(STALLED_KIND) == true }) awaitCancellation()
events.filter { e -> filters.any { it.match(e) } }.forEach(onEach)
if (filters.any { it.kinds?.contains(TRICKLE_KIND) == true }) awaitCancellation()
onEose()
awaitCancellation()
}
override suspend fun count(
ctx: RequestContext,
filters: List<Filter>,
) = events.count { e -> filters.any { it.match(e) } }
override suspend fun submit(
event: Event,
onComplete: (IEventStore.InsertOutcome) -> Unit,
) {
if (events.none { it.id == event.id }) events += event
onComplete(IEventStore.InsertOutcome.Accepted)
}
override suspend fun snapshotIdsForNegentropy(
filters: List<Filter>,
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. */
private class SignedInOnly : PassThroughPolicy() {
@Volatile private var scope: RequestContext? = null
override fun onConnect(
scope: RequestContext,
send: (Message) -> Unit,
) {
this.scope = scope
}
override fun accept(cmd: ReqCmd): PolicyResult<ReqCmd> =
if (scope?.authenticatedUsers?.isNotEmpty() == true) {
PolicyResult.Accepted(cmd)
} else {
PolicyResult.Rejected(MachineReadablePrefix.AUTH_REQUIRED.format("sign in"))
}
}
private class MemoryRelay(
override val backend: MemoryBackend,
signedInOnly: Boolean,
) : RelayServerBase(
policyBuilder = { if (signedInOnly) VerifyPolicy + SignedInOnly() else VerifyPolicy },
parentContext = SupervisorJob(),
negentropySettings = NegentropySettings.Default,
listener = RelayServerListener.None,
limits = RelayLimits(maxMessageLength = 4096),
)
/** One answer as the host would see it. */
private class Recorded : HttpRelayResponse {
var status = 0
val lines = mutableListOf<String>()
var streamed = false
override suspend fun single(
status: Int,
frame: String,
) {
this.status = status
lines += frame
}
override suspend fun stream(lines: suspend HttpRelayLines.() -> Unit) {
status = 200
streamed = true
val out = this.lines
object : HttpRelayLines {
override suspend fun line(frame: String) {
out += frame
}
override suspend fun flush() {}
}.lines()
}
}
private val origin = "https://relay.example"
private val alice = NostrSignerSync()
private val backend = MemoryBackend()
private fun handler(
signedInOnly: Boolean = false,
deadlineMs: Long = 5_000,
) = HttpRelayHandler(MemoryRelay(backend, signedInOnly), origins = { listOf(origin) }, deadlineMs = deadlineMs)
private fun HttpRelayHandler.ask(
command: HttpRelayCommand,
body: String,
authorization: String? = null,
) = Recorded().also { runBlocking { handle(HttpRelayRequest(command, authorization, body.encodeToByteArray()), it) } }
private fun note(
content: String,
at: Long = 1_700_000_000L,
) = alice.sign<Event>(at, 1, emptyArray(), content)
private fun token(
command: HttpRelayCommand,
body: String,
) = alice.sign(HTTPAuthorizationEvent.build(origin + command.path, "POST", body.encodeToByteArray(), System.currentTimeMillis() / 1000) {}).toAuthToken()
@Test
fun aReqStreamsItsEventsAndEndsOnEose() {
val a = note("a")
val b = note("b")
backend.events += listOf(a, b)
val answer = handler().ask(HttpRelayCommand.REQ, """{"kinds":[1]}""")
assertEquals(200, answer.status)
assertTrue(answer.streamed)
assertEquals("""["EOSE","http"]""", answer.lines.last())
assertEquals(
setOf(a.id, b.id),
answer.lines
.dropLast(1)
.map { l -> listOf(a, b).first { it.id in l }.id }
.toSet(),
)
}
@Test
fun anEmptyReqIsOneEoseLine() {
val answer = handler().ask(HttpRelayCommand.REQ, """[{"kinds":[30000]}]""")
assertEquals(200, answer.status)
assertEquals(listOf("""["EOSE","http"]"""), answer.lines)
}
@Test
fun anEventIsAnsweredByItsOkAndAForgeryIsRefused() {
val posted = note("posted")
val ok = handler().ask(HttpRelayCommand.EVENT, posted.toJson())
assertEquals(200, ok.status)
assertEquals(listOf("""["OK","${posted.id}",true,""]"""), ok.lines)
assertTrue(backend.events.any { it.id == posted.id })
val forged = Event(posted.id, posted.pubKey, posted.createdAt, posted.kind, posted.tags, "tampered", posted.sig)
val refused = handler().ask(HttpRelayCommand.EVENT, forged.toJson())
assertEquals(400, refused.status, refused.lines.toString())
assertTrue(refused.lines.single().startsWith("""["OK","${forged.id}",false,"""), refused.lines.toString())
}
@Test
fun aCountIsOneCountLine() {
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())
}
@Test
fun aBodyThatIsNotTheCommandsArgumentsIsA400AndOneOverTheLimitA413() {
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())
}
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)
}
@Test
fun aNip98SignatureSignsTheSessionInAndAnotherSchemeDoesNot() {
backend.events += note("gated")
val gated = handler(signedInOnly = true)
val body = """{"kinds":[1]}"""
val anonymous = gated.ask(HttpRelayCommand.REQ, body)
assertEquals(401, anonymous.status)
assertTrue(anonymous.lines.single().startsWith("""["CLOSED","http","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())
}
@Test
fun aTokenForAnotherBodyOrASecondUseIsRefused() {
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 other = h.ask(HttpRelayCommand.REQ, """{"kinds":[0]}""", token(HttpRelayCommand.REQ, body))
assertEquals(401, other.status)
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<String>()
val need = mutableSetOf<String>()
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","http","""")
.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","http","""), answer.lines.toString())
}
@Test
fun noFirstFrameWithinTheDeadlineIsA503() {
val answer = handler(deadlineMs = 300).ask(HttpRelayCommand.REQ, """{"kinds":[$STALLED_KIND]}""")
assertEquals(503, answer.status)
assertTrue(answer.lines.single().startsWith("""["CLOSED","http","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]}""")
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())
}
@Test
fun aReaderThatStopsReadingIsDroppedAtTheHardStop() {
backend.events += note("found")
val h = HttpRelayHandler(MemoryRelay(backend, false), origins = { listOf(origin) }, deadlineMs = 100, tailGraceMs = 100)
val stalled =
object : HttpRelayResponse {
override suspend fun single(
status: Int,
frame: String,
) = error("streams")
override suspend fun stream(lines: suspend HttpRelayLines.() -> Unit) =
object : HttpRelayLines {
override suspend fun line(frame: String) = awaitCancellation()
override suspend fun flush() {}
}.lines()
}
assertFailsWith<HttpRelayReaderStalled> {
runBlocking { h.handle(HttpRelayRequest(HttpRelayCommand.REQ, null, """{"kinds":[1,$TRICKLE_KIND]}""".encodeToByteArray()), stalled) }
}
}
@Test
fun theStringOnlyConnectStillSendsWireJson() {
backend.events += note("socket")
val out = CopyOnWriteArrayList<String>()
val session = MemoryRelay(backend, false).connect { out += it }
runBlocking { session.receive("""["REQ","s",{"kinds":[1]}]""") }
assertTrue(out.any { it.startsWith("""["EVENT","s",""") } && out.contains("""["EOSE","s"]"""), out.toString())
session.close()
}
private companion object {
const val STALLED_KIND = 7
const val TRICKLE_KIND = 8
}
}