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..cb5a865ea6 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,20 @@ */ 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. + * 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. */ enum class HttpRelayCommand( val path: String, @@ -56,32 +41,33 @@ enum class HttpRelayCommand( 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. + * 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 parse(body: String): Parsed? { - val tree = - try { - Json.parseToJsonElement(body) - } catch (_: SerializationException) { - return null + fun frameOf(body: String): String? { + val text = body.trim() + return when (this) { + REQ, COUNT -> { + val filters = + when { + text.startsWith('{') -> text + text.startsWith('[') && text.endsWith(']') -> text.substring(1, text.length - 1).trim().ifEmpty { return null } + else -> return null + } + "[\"${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 (!text.startsWith('{')) return null + "[\"${EventCmd.LABEL}\",$text]" + } + } + } /** Whether [message] is the last frame of this command's answer. */ fun ends(message: Message): Boolean = @@ -90,48 +76,13 @@ 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 } - 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" fun forPath(path: String): HttpRelayCommand? = entries.firstOrNull { it.path == path } 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..a4054bbd08 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt @@ -21,7 +21,6 @@ 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 @@ -29,6 +28,7 @@ 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.CancellationException import kotlinx.coroutines.CompletableDeferred import kotlinx.coroutines.TimeoutCancellationException import kotlinx.coroutines.channels.Channel @@ -36,22 +36,27 @@ 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 -/** 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 +65,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 +78,146 @@ 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, /** 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 +227,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 } } @@ -252,8 +284,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, once. + * 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 { @@ -262,34 +294,46 @@ 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") - 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 + if (accepted.isEmpty()) return Proof.Refused(MachineReadablePrefix.AUTH_REQUIRED.format("this relay names no url to sign")) + // 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}")) } - /** 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() + 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 } } + +/** 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") + +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..370ca0395e 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,16 @@ 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.nip77Negentropy.NegErrMessage +import com.vitorpamplona.quartz.nip01Core.store.RejectionReason /** - * 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 +42,24 @@ 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..31695d8da3 100644 --- a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt @@ -20,25 +20,27 @@ */ 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.core.HexKey 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 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 kotlinx.coroutines.SupervisorJob @@ -49,8 +51,11 @@ 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. */ +/** 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 { @@ -81,11 +86,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. */ @@ -108,15 +108,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 +156,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 +183,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 +197,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 +219,25 @@ 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,89 +251,48 @@ 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 - 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()) } - @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","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]}""") + 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 +322,131 @@ 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 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 submit( + event: Event, + onComplete: (IEventStore.InsertOutcome) -> Unit, + ): Unit = error("db is down") + } + 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()) + } + + @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 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("", "[]", "[ ]", "[{}", "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