mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 11:18:24 +00:00
Merge pull request #4214 from vitorpamplona/claude/nip-fe-http-followup
fix(quartz): NIP-FE review fixes, NIP-98 sign-in policy vote, frames without sub id, drop /neg
This commit is contained in:
+5
-16
@@ -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<HexKey> = 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<HexKey>,
|
||||
incoming: suspend (RelaySession) -> Unit,
|
||||
) {
|
||||
val session = connect(sink, authenticatedUsers)
|
||||
val session = connect(sink)
|
||||
try {
|
||||
incoming(session)
|
||||
} finally {
|
||||
|
||||
+19
-9
@@ -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<HexKey> = 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<HexKey> = initialAuthenticatedUsers.toSet()
|
||||
private var authenticatedUsers = setOf<HexKey>()
|
||||
|
||||
/**
|
||||
* 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)
|
||||
|
||||
+9
@@ -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<EventCmd> =
|
||||
if (isAuthenticated()) {
|
||||
PolicyResult.Accepted(cmd)
|
||||
|
||||
+15
@@ -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<T : Command> {
|
||||
val reason: String,
|
||||
) : PolicyResult<T>
|
||||
}
|
||||
|
||||
/** 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"
|
||||
|
||||
+4
@@ -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<AuthCmd> = 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
|
||||
}
|
||||
|
||||
+4
@@ -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(
|
||||
|
||||
+4
@@ -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
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
+28
-77
@@ -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<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. */
|
||||
/**
|
||||
* 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 }
|
||||
|
||||
+142
-98
@@ -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<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(),
|
||||
/** 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<HexKey>,
|
||||
signedIn: 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,
|
||||
) {
|
||||
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<Frame>,
|
||||
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)
|
||||
}
|
||||
|
||||
+15
-15
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
+166
-75
@@ -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<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. */
|
||||
@@ -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<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]}""")
|
||||
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<HttpRelayReaderStalled> {
|
||||
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
|
||||
|
||||
Reference in New Issue
Block a user