diff --git a/geode/README.md b/geode/README.md index 4c45f710fc..11732aaad9 100644 --- a/geode/README.md +++ b/geode/README.md @@ -3,8 +3,8 @@ A standalone [Nostr](https://github.com/nostr-protocol/nips) relay for the JVM, built on Quartz's relay-server code (Ktor CIO). It speaks the core relay protocol plus NIP-11 (info doc), NIP-42 (AUTH), NIP-45 (COUNT), NIP-50 -(full-text search), NIP-77 (Negentropy sync), and NIP-86 (relay management), -stores events in SQLite (or a filesystem backend), and can mirror upstream +(full-text search), NIP-77 (Negentropy sync), NIP-86 (relay management) and +NIP-FE (REQ/COUNT/EVENT over plain HTTP), stores events in SQLite (or a filesystem backend), and can mirror upstream relays strfry-router style. geode depends only on `:quartz` — no Android, no Compose. `amy serve` (the @@ -97,8 +97,34 @@ geode --version Key sections: `[info]` (NIP-11 doc), `[network]` (bind + thread pools), `[database]` (SQLite path/tuning), `[options]` (AUTH / verify / search), -`[authorization]` (allow/deny lists), `[[mirror]]` (upstream mirroring), and -`[admin]` (NIP-86 management). See the example file for every knob. +`[authorization]` (allow/deny lists), `[[mirror]]` (upstream mirroring), +`[http]` (NIP-FE limits) and `[admin]` (NIP-86 management). See the example +file for every knob. + +### Commands over HTTP (NIP-FE) + +Besides the websocket, geode answers one command per HTTP `POST` to the relay +URL: the body is the frame you would send on the socket, and the answer is the +socket's frames as NDJSON, streamed — no socket, no subscription left open: + +```bash +curl -N --compressed -d '["REQ","q",{"kinds":[1],"limit":2}]' http://localhost:7447/ +# ["EVENT","q",{"id":"…","kind":1,…}] +# ["EVENT","q",{"id":"…","kind":1,…}] +# ["EOSE","q"] +curl -d '["COUNT","c",{"kinds":[1]}]' http://localhost:7447/ # ["COUNT","c",{"count":2}] +curl -d "[\"EVENT\",$(cat signed-event.json)]" http://localhost:7447/ # ["OK","",true,""] +``` + +A body that does not end on `EOSE`/`CLOSED` (REQ), `COUNT`/`CLOSED` (COUNT) or +`OK` (EVENT) was cut off. On an AUTH-gated relay the answer is `401` until the +request carries a NIP-98 `Authorization: Nostr …` header whose `u` is the +relay's http URL and whose `payload` is the body's sha256. NIP-86 admin calls +share the URL, told apart by `Content-Type: application/nostr+json+rpc`. +Streamed answers are gzipped for clients that send `Accept-Encoding: gzip` +(`curl --compressed`), sync-flushed so events still arrive as they are found; +`[http].gzip = false` turns it off. Quartz's `HttpRelayClient` does all of this +for clients, over OkHttp via `OkHttpRelayTransport` or any `HttpRelayTransport`. ## Verbs diff --git a/geode/config.example.toml b/geode/config.example.toml index 4651105623..89f6a6c196 100644 --- a/geode/config.example.toml +++ b/geode/config.example.toml @@ -19,7 +19,8 @@ contact = "admin@example.com" # pubkey = "..." # Override the supported NIPs advertised on the NIP-11 endpoint. If # omitted, the relay advertises the NIPs it actually implements. -# supported_nips = [1, 9, 11, 40, 42, 45, 50, 62] +# Hex-named NIPs go in as strings. +# supported_nips = [1, 9, 11, 40, 42, 45, 50, 62, "FE"] [network] host = "0.0.0.0" @@ -165,6 +166,43 @@ require_auth = false # backfill_seconds = 3600 # filter = '{"kinds":[0,1,3,7],"#t":["nostr"]}' +[http] +# NIP-FE: relay commands over HTTP. One REQ, COUNT or EVENT frame per +# POST to the relay URL, exactly as it would go on the websocket, +# answered as NDJSON (application/x-ndjson) in the socket's own frames and +# streamed as it is found; nothing stays open afterwards. NIP-86 calls +# share the URL, told apart by their application/nostr+json+rpc type. +# NIP-98 `Authorization: Nostr ...` headers sign a request in, exactly as +# NIP-42 AUTH would on the socket. On by default; turning it off also +# drops "FE" from the default NIP-11 list. +enabled = true +# Every request is its own connection, so the websocket's +# per-connection limits don't bound HTTP clients. These do: over the +# per-client cap a request gets 429, over the global cap 503, both with +# Retry-After. 0 = no limit. +max_concurrent_requests = 256 +max_requests_per_client = 16 +# A body still arriving after this is dropped with 408, freeing its slot. +body_timeout_seconds = 10 +# An answer still running after this ends on a CLOSED line. +deadline_seconds = 30 +max_body_bytes = 524288 +# Gzip streamed answers for clients that send Accept-Encoding: gzip. Each +# batch of lines is sync-flushed, so events still arrive as they are +# found. Turn off if a proxy in front already compresses. +gzip = true +retry_after_seconds = 1 +# Other addresses this relay is reachable at: a NIP-98 token may name +# any of them, as well as [info].relay_url. +# alternate_urls = ["ws://youraddress.onion/"] +# Behind a reverse proxy every request comes from the proxy's address. +# List the proxies here and the per-client cap counts the address the +# proxy writes in client_address_header instead (its last entry). The +# header is ignored from anyone else. Also set `proxy_buffering off` +# (nginx) or equivalent; the relay sends X-Accel-Buffering: no. +# trusted_proxies = ["127.0.0.1"] +# client_address_header = "X-Forwarded-For" + [admin] # NIP-86 relay management API. When `pubkeys` is non-empty, the relay # accepts HTTP POST application/nostr+json+rpc on the same URL, diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/KtorRelay.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/KtorRelay.kt index 0a1af895c1..4385521331 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/KtorRelay.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/KtorRelay.kt @@ -20,13 +20,16 @@ */ package com.vitorpamplona.geode +import com.vitorpamplona.geode.server.HttpCommandSettings import com.vitorpamplona.geode.server.Nip11HttpRoute import com.vitorpamplona.geode.server.Nip86HttpRoute +import com.vitorpamplona.geode.server.NipFEHttpRoute import com.vitorpamplona.geode.server.WebSocketSessionPump import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.NoticeMessage import com.vitorpamplona.quartz.nip01Core.relay.normalizer.toHttp import com.vitorpamplona.quartz.nip01Core.relay.server.RelaySession import com.vitorpamplona.quartz.nip86RelayManagement.server.Nip86HttpHandler +import com.vitorpamplona.quartz.nipFERelayOverHttp.HttpRelayHandler import io.ktor.server.application.install import io.ktor.server.application.serverConfig import io.ktor.server.cio.CIO @@ -34,6 +37,7 @@ import io.ktor.server.cio.CIOApplicationEngine import io.ktor.server.engine.connector import io.ktor.server.engine.embeddedServer import io.ktor.server.routing.get +import io.ktor.server.routing.options import io.ktor.server.routing.post import io.ktor.server.routing.routing import io.ktor.server.websocket.WebSockets @@ -77,6 +81,11 @@ class KtorRelay( val workerGroupSize: Int? = null, /** Ktor CIO call-handling thread count. `null` keeps Ktor's default. */ val callGroupSize: Int? = null, + /** + * NIP-FE: relay commands POSTed to the relay's URL, beside NIP-86. On by default; null turns them + * off, every POST going to NIP-86 again (and the operator should then drop `FE` from NIP-11). + */ + val httpCommands: HttpCommandSettings? = HttpCommandSettings(), ) { /** * NIP-86 HTTP adapter. Wraps the engine's [RelayEngine.nip86Server] @@ -104,6 +113,20 @@ class KtorRelay( private val nip11Route = Nip11HttpRoute(liveJson = { relay.info.json }) + /** + * NIP-FE. Each request runs on its own session of the same engine, so the websocket's policies + * and limits apply. A NIP-98 token must name `relay.url` read as http(s), or one of the + * configured alternate URLs (a .onion). + */ + private val nipFERoute = + httpCommands?.let { settings -> + val origins = (listOf(relay.url) + settings.alternateUrls).map { it.toHttp() } + NipFEHttpRoute( + handler = HttpRelayHandler(relay.server, origins = { origins }, deadline = settings.deadline), + settings = settings, + ) + } + private var engine: CIOApplicationEngine? = null private var resolvedPort: Int = -1 @@ -164,12 +187,18 @@ class KtorRelay( get(path) { nip11Route.handle(call) } - // NIP-86: POST application/nostr+json+rpc with a NIP-98 - // signed Authorization header → JSON-RPC dispatch. - // Always mounted; an empty admin allow-list on the engine just means - // every request fails the allow-list check (403). + // Two POSTs share the relay URL, told apart by Content-Type: + // - NIP-86: application/nostr+json+rpc with a NIP-98 signed + // Authorization header → JSON-RPC dispatch. Always mounted; an + // empty admin allow-list just means every call fails it (403). + // - NIP-FE: anything else is one REQ/COUNT/EVENT frame, answered + // as NDJSON in the socket's own frames. post(path) { - nip86Route.handle(call) + val commands = nipFERoute + if (commands != null && commands.isCommand(call)) commands.handle(call) else nip86Route.handle(call) + } + nipFERoute?.let { route -> + options(path) { route.preflight(call) } } webSocket(path) { if (shuttingDown) { diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt index 1ed17bbcdb..e4ddc8bdea 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt @@ -35,9 +35,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl import com.vitorpamplona.quartz.nip01Core.relay.normalizer.displayUrl import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl import com.vitorpamplona.quartz.nip01Core.relay.server.policies.EmptyPolicy -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.OptionalAuthPolicy import com.vitorpamplona.quartz.nip01Core.relay.server.policies.RejectFutureEventsPolicy import com.vitorpamplona.quartz.nip01Core.relay.server.policies.VerifyAuthOnlyPolicy import com.vitorpamplona.quartz.nip01Core.relay.server.policies.VerifyPolicy @@ -356,6 +354,7 @@ private fun serve(args: Array) { connectionGroupSize = config.network.connection_group_size, workerGroupSize = config.network.worker_group_size, callGroupSize = config.network.call_group_size, + httpCommands = config.http.toSettings(), ).start() // `[[mirror]]` upstreams: dial each configured relay and stream its @@ -526,9 +525,9 @@ private fun composePolicy( val pieces = mutableListOf() if (requireAuth) { - pieces += FullAuthPolicy(advertisedUrl) + pieces += SignInPolicy(advertisedUrl) } else if (optionalAuth) { - pieces += OptionalAuthPolicy(advertisedUrl) + pieces += OptionalSignInPolicy(advertisedUrl) } config.options.reject_future_seconds?.let { secs -> diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/RelayInfo.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/RelayInfo.kt index 4170fa14b4..92f294b519 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/RelayInfo.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/RelayInfo.kt @@ -62,9 +62,10 @@ data class RelayInfo( * - 62 NIP-62 right to vanish * - 77 NIP-77 negentropy reconciliation * - 86 NIP-86 relay management API (when admin pubkeys configured) + * - FE NIP-FE relay commands over HTTP (KtorRelay; `[http].enabled`) */ val SUPPORTED_NIPS: List = - listOf("1", "9", "11", "40", "42", "45", "50", "62", "77", "86") + listOf("1", "9", "11", "40", "42", "45", "50", "62", "77", "86", "FE") /** Pre-built default for `RelayEngine(url = ...)` — advertises the supported NIPs. */ fun default(url: NormalizedRelayUrl): RelayInfo = diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/SignInPolicies.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/SignInPolicies.kt new file mode 100644 index 0000000000..b246219bab --- /dev/null +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/SignInPolicies.kt @@ -0,0 +1,43 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.geode + +import com.vitorpamplona.quartz.nip01Core.core.HexKey +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import com.vitorpamplona.quartz.nip01Core.relay.server.policies.FullAuthPolicy +import com.vitorpamplona.quartz.nip01Core.relay.server.policies.OptionalAuthPolicy + +/** + * Geode's AUTH decides nothing past the NIP-42 proof, so a key a NIP-98 header proved on a NIP-FE + * command signs in just the same: the transport's proof is all the socket's would have been. + */ +internal class SignInPolicy( + relay: NormalizedRelayUrl, +) : FullAuthPolicy(relay) { + override suspend fun authorizeTransport(pubkey: HexKey): String? = null +} + +/** [SignInPolicy] for optional AUTH. */ +internal class OptionalSignInPolicy( + relay: NormalizedRelayUrl, +) : OptionalAuthPolicy(relay) { + override suspend fun authorizeTransport(pubkey: HexKey): String? = null +} diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/config/StaticConfig.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/config/StaticConfig.kt index 8cfecb4e41..e12781ad13 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/config/StaticConfig.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/config/StaticConfig.kt @@ -21,12 +21,15 @@ package com.vitorpamplona.geode.config import cc.ekblad.toml.decode +import cc.ekblad.toml.model.TomlValue import cc.ekblad.toml.tomlMapper import com.vitorpamplona.geode.RelayInfo +import com.vitorpamplona.geode.server.HttpCommandSettings import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl import com.vitorpamplona.quartz.nip11RelayInfo.Nip11RelayInformation import java.io.File +import kotlin.time.Duration.Companion.seconds /** * Operator-facing **boot-time** configuration. Parsed once from a TOML @@ -45,10 +48,14 @@ data class StaticConfig( val authorization: AuthorizationSection = AuthorizationSection(), val admin: AdminSection = AdminSection(), val negentropy: NegentropySection = NegentropySection(), + val http: HttpSection = HttpSection(), /** `[[mirror]]` entries — upstream relays this relay streams from. */ val mirror: List = emptyList(), ) { - fun resolveInfo(fullTextSearch: Boolean = true): RelayInfo = + fun resolveInfo( + fullTextSearch: Boolean = true, + httpCommands: Boolean = http.enabled, + ): RelayInfo = RelayInfo( Nip11RelayInformation( name = info.name ?: RelayInfo.NAME, @@ -59,10 +66,10 @@ data class StaticConfig( software = info.software ?: RelayInfo.SOFTWARE, version = info.version ?: RelayInfo.VERSION, supported_nips = - info.supported_nips?.map(Int::toString) + info.supported_nips // An explicit [info] nips list is operator-authoritative; - // the default list stays honest about search. - ?: if (fullTextSearch) RelayInfo.SUPPORTED_NIPS else RelayInfo.SUPPORTED_NIPS - "50", + // the default list stays honest about search and HTTP commands. + ?: RelayInfo.SUPPORTED_NIPS.filter { (fullTextSearch || it != "50") && (httpCommands || it != "FE") }, privacy_policy = info.privacy_policy, terms_of_service = info.terms_of_service, relay_countries = info.relay_countries, @@ -80,7 +87,8 @@ data class StaticConfig( val icon: String? = null, val software: String? = null, val version: String? = null, - val supported_nips: List? = null, + /** NIP numbers, and the hex-named NIPs as strings: `[1, 11, 42, "FE"]`. */ + val supported_nips: List? = null, val privacy_policy: String? = null, val terms_of_service: String? = null, val relay_countries: List? = null, @@ -223,6 +231,59 @@ data class StaticConfig( val live_index: Boolean = true, ) + /** + * NIP-FE: relay commands over HTTP — one REQ, COUNT or EVENT frame POSTed to the relay's URL, + * answered as NDJSON in the socket's own frames. Each request is its own connection, so the + * concurrency caps here stand in for the websocket's per-connection limits. + */ + data class HttpSection( + val enabled: Boolean = true, + /** Requests running at once across all clients; over it, 503. 0 = no limit. */ + val max_concurrent_requests: Int = 256, + /** Requests one client address may run at once; over it, 429. 0 = no limit. */ + val max_requests_per_client: Int = 16, + /** How long a request body may take to arrive before the request is dropped with 408. */ + val body_timeout_seconds: Long = 10, + /** How long one answer may run before it ends on a `CLOSED` line. */ + val deadline_seconds: Long = 30, + /** Largest request body read; larger is 413. */ + val max_body_bytes: Int = 512 * 1024, + /** Gzip streamed answers when the client sends `Accept-Encoding: gzip`, flushed line by line. */ + val gzip: Boolean = true, + /** `Retry-After` on the relay's own 429 and 503. */ + val retry_after_seconds: Int = 1, + /** + * Other `ws(s)://` URLs this relay is reachable at (a `.onion` beside the clearnet name). + * A NIP-98 token's `u` may name any of them, read as http(s), or `[info].relay_url`. + */ + val alternate_urls: List = emptyList(), + /** + * Addresses of the reverse proxies in front of the relay. Only from these peers is + * [client_address_header] believed when counting a client's requests. + */ + val trusted_proxies: List = emptyList(), + val client_address_header: String = "X-Forwarded-For", + ) { + /** The transport settings, or null when commands over HTTP are off. */ + fun toSettings(): HttpCommandSettings? = + if (!enabled) { + null + } else { + HttpCommandSettings( + maxConcurrent = max_concurrent_requests, + maxPerClient = max_requests_per_client, + bodyTimeout = body_timeout_seconds.seconds, + deadline = deadline_seconds.seconds, + maxBodyBytes = max_body_bytes, + compress = gzip, + retryAfterSeconds = retry_after_seconds, + alternateUrls = alternate_urls.map { it.normalizeRelayUrl() }, + trustedProxies = trusted_proxies.toSet(), + clientAddressHeader = client_address_header, + ) + } + } + data class AuthorizationSection( val pubkey_whitelist: List = emptyList(), val pubkey_blacklist: List = emptyList(), @@ -297,6 +358,12 @@ data class StaticConfig( * the store. */ fun validate() { + require(http.max_concurrent_requests >= 0) { "[http].max_concurrent_requests must be >= 0 (0 = no limit), got ${http.max_concurrent_requests}" } + require(http.max_requests_per_client >= 0) { "[http].max_requests_per_client must be >= 0 (0 = no limit), got ${http.max_requests_per_client}" } + require(http.body_timeout_seconds > 0) { "[http].body_timeout_seconds must be > 0, got ${http.body_timeout_seconds}" } + require(http.deadline_seconds > 0) { "[http].deadline_seconds must be > 0, got ${http.deadline_seconds}" } + require(http.max_body_bytes > 0) { "[http].max_body_bytes must be > 0, got ${http.max_body_bytes}" } + require(http.retry_after_seconds >= 0) { "[http].retry_after_seconds must be >= 0, got ${http.retry_after_seconds}" } database.readers?.let { require(it >= 1) { "[database].readers must be >= 1 (got $it); a 0/negative pool can never answer a query" } } @@ -308,7 +375,11 @@ data class StaticConfig( } companion object { - private val mapper = tomlMapper { } + private val mapper = + tomlMapper { + // `supported_nips = [1, 11, "FE"]`: numbered NIPs are integers, hex-named ones strings. + decoder { it: TomlValue.Integer -> it.value.toString() } + } fun fromToml(toml: String): StaticConfig = mapper.decode(toml) diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/server/BoundedBody.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/server/BoundedBody.kt new file mode 100644 index 0000000000..571c6312b2 --- /dev/null +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/server/BoundedBody.kt @@ -0,0 +1,54 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.geode.server + +import io.ktor.http.HttpHeaders +import io.ktor.server.application.ApplicationCall +import io.ktor.server.request.receiveChannel +import io.ktor.utils.io.readAvailable + +/** + * Reads the request body up to [cap] bytes. Returns null when it is larger — by its declared + * `Content-Length` or by what actually arrives — without reading past the cap; the caller answers 413. + * The buffer is sized to the declared length, or grows from a small start, so a 50-byte command does + * not cost a cap-sized allocation. + */ +internal suspend fun readBoundedBody( + call: ApplicationCall, + cap: Int, +): ByteArray? { + val declared = call.request.headers[HttpHeaders.ContentLength]?.toLongOrNull() + if (declared != null && declared > cap) return null + val ch = call.receiveChannel() + // One byte past the declared length, so a body longer than it claimed is still caught at the cap. + var buf = ByteArray(if (declared != null) declared.toInt() + 1 else minOf(cap + 1, INITIAL_BODY_BUFFER)) + var pos = 0 + while (pos <= cap) { + if (pos == buf.size) buf = buf.copyOf(minOf(cap + 1, buf.size * 2)) + val read = ch.readAvailable(buf, pos, buf.size - pos) + if (read <= 0) break + pos += read + } + if (pos > cap) return null + return if (pos == buf.size) buf else buf.copyOf(pos) +} + +private const val INITIAL_BODY_BUFFER = 4 * 1024 diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/server/GzipLines.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/server/GzipLines.kt new file mode 100644 index 0000000000..f1b50e171a --- /dev/null +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/server/GzipLines.kt @@ -0,0 +1,88 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.geode.server + +import io.ktor.utils.io.ByteWriteChannel +import io.ktor.utils.io.writeFully +import java.io.ByteArrayOutputStream +import java.util.zip.Deflater +import java.util.zip.GZIPOutputStream + +/** + * A gzip body written line by line onto [out], for NIP-FE's streamed answers. Every [flush] is a + * gzip sync flush: the deflater gives up everything it holds, byte-aligned, so the client can + * inflate every line written so far. Without it the compressor sits on the first lines until its + * window fills, and streaming delivers nothing until the answer is nearly done. + */ +internal class GzipLines( + private val out: ByteWriteChannel, +) { + private val compressed = Pending() + + // Level 1: on Nostr events it compresses to ~42% of raw against ~39% at the JDK's default 6, for + // about 2.5x less CPU; ids, keys and signatures are hex and barely compress at any level. + private val gzip = + object : GZIPOutputStream(compressed, BUFFER, true) { + init { + def.setLevel(Deflater.BEST_SPEED) + } + } + + /** Compresses [line] and its newline. Moves compressed bytes to the socket once enough piled up, so a long burst still meets backpressure. */ + suspend fun line(line: String) { + gzip.write(line.encodeToByteArray()) + gzip.write(NEWLINE) + if (compressed.size() >= BUFFER) drain() + } + + /** A sync flush, then everything onto the socket. */ + suspend fun flush() { + gzip.flush() + drain() + out.flush() + } + + /** Writes the gzip trailer; the body is complete after this. */ + suspend fun finish() { + gzip.finish() + drain() + out.flush() + } + + /** Frees the deflater, finished or not. */ + fun close() = gzip.close() + + private suspend fun drain() { + if (compressed.size() == 0) return + compressed.writeTo(out) + compressed.reset() + } + + /** The compressed bytes not yet on the socket, written from its own array rather than a copy. */ + private class Pending : ByteArrayOutputStream(BUFFER) { + suspend fun writeTo(out: ByteWriteChannel) = out.writeFully(buf, 0, count) + } + + companion object { + private const val BUFFER = 16 * 1024 + private val NEWLINE = byteArrayOf('\n'.code.toByte()) + } +} diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/server/HttpAdmission.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/server/HttpAdmission.kt new file mode 100644 index 0000000000..939d3ab22d --- /dev/null +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/server/HttpAdmission.kt @@ -0,0 +1,89 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.geode.server + +import java.util.concurrent.ConcurrentHashMap +import java.util.concurrent.atomic.AtomicInteger + +/** + * How many NIP-FE requests run at once, per client and in all. Each HTTP request is its own + * connection, so the per-connection limits the websocket leans on (open subscriptions, one search at + * a time) bound nothing here; this does. A request over its client's share is that client's to + * slow down (429); one over the relay's is everyone's (503). + */ +internal class HttpAdmission( + /** Requests running at once across all clients; 0 is no limit. */ + private val maxConcurrent: Int, + /** Requests one client may run at once; 0 is no limit. */ + private val maxPerClient: Int, +) { + enum class Verdict { ADMITTED, CLIENT_BUSY, RELAY_BUSY } + + private val running = AtomicInteger() + private val perClient = ConcurrentHashMap() + + /** Number of requests running now. */ + val inFlight: Int get() = running.get() + + /** Runs [block] if [client] and the relay have room; otherwise says which of them had none. */ + suspend fun admit( + client: String, + block: suspend () -> Unit, + ): Verdict { + val verdict = enter(client) + if (verdict != Verdict.ADMITTED) return verdict + try { + block() + } finally { + leave(client) + } + return verdict + } + + private fun enter(client: String): Verdict { + var clientFull = false + perClient.compute(client) { _, n -> + val now = n ?: 0 + if (maxPerClient in 1..now) { + clientFull = true + n + } else { + now + 1 + } + } + if (clientFull) return Verdict.CLIENT_BUSY + if (running.incrementAndGet().let { maxConcurrent in 1 until it }) { + running.decrementAndGet() + release(client) + return Verdict.RELAY_BUSY + } + return Verdict.ADMITTED + } + + private fun leave(client: String) { + running.decrementAndGet() + release(client) + } + + private fun release(client: String) { + perClient.computeIfPresent(client) { _, n -> if (n <= 1) null else n - 1 } + } +} diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/server/HttpCommandSettings.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/server/HttpCommandSettings.kt new file mode 100644 index 0000000000..9d4e75d2d4 --- /dev/null +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/server/HttpCommandSettings.kt @@ -0,0 +1,58 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.geode.server + +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import com.vitorpamplona.quartz.nipFERelayOverHttp.HttpRelayHandler +import kotlin.time.Duration +import kotlin.time.Duration.Companion.seconds + +/** + * NIP-FE (relay commands over HTTP) as [com.vitorpamplona.geode.KtorRelay] serves it: a `POST` to + * the relay's URL carrying one REQ, COUNT or EVENT frame. The engine's policies and limits apply as + * they do on the websocket; these bound what the websocket's per-connection limits cannot, since + * every request is its own connection. + */ +data class HttpCommandSettings( + /** Requests running at once across all clients before the rest get 503; 0 is no limit. */ + val maxConcurrent: Int = 256, + /** Requests one client address may run at once before its next gets 429; 0 is no limit. */ + val maxPerClient: Int = 16, + /** How long the body may take to arrive; past it, 408. It holds an admission slot meanwhile. */ + val bodyTimeout: Duration = 10.seconds, + /** How long one answer may run, first byte to last. */ + val deadline: Duration = HttpRelayHandler.DEFAULT_DEADLINE, + /** The largest body read; the engine's own message limit, when it has one and it is smaller, wins. */ + val maxBodyBytes: Int = 512 * 1024, + /** Gzip streamed answers for clients that accept it, sync-flushed so lines still arrive as found. */ + val compress: Boolean = true, + /** The `Retry-After` sent with a 429 or 503 the relay decides itself. */ + val retryAfterSeconds: Int = 1, + /** Other URLs this relay answers at (its .onion); a NIP-98 `u` may name any of them. */ + val alternateUrls: List = emptyList(), + /** + * Peers whose [clientAddressHeader] is believed: the reverse proxies in front of the relay. A + * request from anyone else is counted under its own address, whatever the header says. + */ + val trustedProxies: Set = emptySet(), + /** Where a trusted proxy writes the client's address; its last entry is the one the proxy saw. */ + val clientAddressHeader: String = "X-Forwarded-For", +) diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/server/Nip86HttpRoute.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/server/Nip86HttpRoute.kt index 5c3e49d800..13dd281efb 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/server/Nip86HttpRoute.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/server/Nip86HttpRoute.kt @@ -26,9 +26,7 @@ import io.ktor.http.HttpHeaders import io.ktor.http.HttpStatusCode import io.ktor.server.application.ApplicationCall import io.ktor.server.request.header -import io.ktor.server.request.receiveChannel import io.ktor.server.response.respondText -import io.ktor.utils.io.readAvailable /** * Ktor adapter for the canonical NIP-86 HTTP flow encapsulated by @@ -55,7 +53,12 @@ internal class Nip86HttpRoute( private val handler: Nip86HttpHandler, ) { suspend fun handle(call: ApplicationCall) { - val body = readBoundedBody(call, handler.maxBodyBytes) ?: return // 413 already sent + val body = + readBoundedBody(call, handler.maxBodyBytes) ?: return call.respondText( + "request body exceeds ${handler.maxBodyBytes}-byte cap", + ContentType.Text.Plain, + HttpStatusCode.PayloadTooLarge, + ) val authHeader = call.request.header(HttpHeaders.Authorization) when (val r = handler.handle(authHeader, body)) { @@ -111,43 +114,6 @@ internal class Nip86HttpRoute( } } - /** - * Bounded read using [cap]. Returns null after sending a 413 if - * the request body exceeds the cap — either the declared - * `Content-Length` or what we actually pull off the wire. - */ - private suspend fun readBoundedBody( - call: ApplicationCall, - cap: Int, - ): ByteArray? { - val declared = call.request.headers[HttpHeaders.ContentLength]?.toLongOrNull() - if (declared != null && declared > cap) { - call.respondText( - "request body exceeds $cap-byte cap", - ContentType.Text.Plain, - HttpStatusCode.PayloadTooLarge, - ) - return null - } - val ch = call.receiveChannel() - val buf = ByteArray(cap + 1) - var pos = 0 - while (pos <= cap) { - val read = ch.readAvailable(buf, pos, buf.size - pos) - if (read <= 0) break - pos += read - } - if (pos > cap) { - call.respondText( - "request body exceeds $cap-byte cap", - ContentType.Text.Plain, - HttpStatusCode.PayloadTooLarge, - ) - return null - } - return buf.copyOfRange(0, pos) - } - /** * Single-line stderr audit. Operators grep on "nip86" / pubkey / * method without needing a logging framework. Best-effort — a diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/server/NipFEHttpRoute.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/server/NipFEHttpRoute.kt new file mode 100644 index 0000000000..57aeefbfa4 --- /dev/null +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/server/NipFEHttpRoute.kt @@ -0,0 +1,231 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.geode.server + +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.MachineReadablePrefix +import com.vitorpamplona.quartz.nip86RelayManagement.server.Nip86HttpHandler +import com.vitorpamplona.quartz.nipFERelayOverHttp.HttpRelayHandler +import com.vitorpamplona.quartz.nipFERelayOverHttp.HttpRelayLines +import com.vitorpamplona.quartz.nipFERelayOverHttp.HttpRelayRequest +import com.vitorpamplona.quartz.nipFERelayOverHttp.HttpRelayResponse +import com.vitorpamplona.quartz.nipFERelayOverHttp.HttpRelayStatus +import io.ktor.http.ContentType +import io.ktor.http.HttpHeaders +import io.ktor.http.HttpStatusCode +import io.ktor.server.application.ApplicationCall +import io.ktor.server.request.header +import io.ktor.server.response.header +import io.ktor.server.response.respond +import io.ktor.server.response.respondBytesWriter +import io.ktor.server.response.respondText +import io.ktor.utils.io.writeStringUtf8 +import kotlinx.coroutines.TimeoutCancellationException +import kotlinx.coroutines.withTimeout + +/** + * NIP-FE over Ktor: the host half of [HttpRelayHandler], on POSTs to the relay's URL that are not + * NIP-86 calls (see [isCommand]). It admits the request, reads the body up to + * the cap, and writes the handler's answer as `application/x-ndjson` — one refusal line with its + * status, or a 200 streamed and flushed frame by frame — with the headers the status calls for + * (`WWW-Authenticate` on 401, `Retry-After` on 429 and 503) and the CORS and no-buffering headers + * every answer carries. + */ +internal class NipFEHttpRoute( + private val handler: HttpRelayHandler, + private val settings: HttpCommandSettings, +) { + private val admission = HttpAdmission(settings.maxConcurrent, settings.maxPerClient) + + private val bodyCap: Int = minOf(settings.maxBodyBytes.toLong(), handler.maxBodyBytes ?: Long.MAX_VALUE).toInt() + + /** Requests being answered right now. */ + val inFlight: Int get() = admission.inFlight + + /** A CORS preflight: any origin may POST with a NIP-98 `Authorization`; no cookies are involved. */ + suspend fun preflight(call: ApplicationCall) { + call.response.header(HttpHeaders.AccessControlAllowOrigin, "*") + call.response.header(HttpHeaders.AccessControlAllowMethods, "POST, OPTIONS") + call.response.header(HttpHeaders.AccessControlAllowHeaders, "Authorization, Content-Type") + call.response.header(HttpHeaders.AccessControlMaxAge, PREFLIGHT_MAX_AGE_SECONDS.toString()) + call.respond(HttpStatusCode.NoContent) + } + + /** + * Whether a POST to the relay's URL is a NIP-FE command: anything but NIP-86's + * `application/nostr+json+rpc`, since commands need no `Content-Type` at all. + */ + fun isCommand(call: ApplicationCall): Boolean { + // Compared as text: parsing it would throw on a malformed header, and a command needs none. + val type = call.request.header(HttpHeaders.ContentType) ?: return true + return !type.substringBefore(';').trim().equals(Nip86HttpHandler.CONTENT_TYPE, ignoreCase = true) + } + + /** Answers the command in the body. Admission runs first, so a refused request spends no NIP-98 token. */ + suspend fun handle(call: ApplicationCall) { + call.response.header(HttpHeaders.AccessControlAllowOrigin, "*") + call.response.header(HttpHeaders.AccessControlExposeHeaders, "${HttpHeaders.WWWAuthenticate}, ${HttpHeaders.RetryAfter}") + call.response.header(HttpHeaders.CacheControl, "no-store") + // nginx and friends buffer responses by default, which holds back the lines streaming delivers. + call.response.header(ACCEL_BUFFERING, "no") + + val verdict = + admission.admit(clientOf(call)) { + // Bounded in time too: a body trickled in forever would hold this admission slot forever. + val body = + try { + withTimeout(settings.bodyTimeout) { readBoundedBody(call, bodyCap) } + } catch (_: TimeoutCancellationException) { + call.response.header(HttpHeaders.Connection, "close") + return@admit respondLine( + call, + REQUEST_TIMEOUT, + HttpRelayHandler.notice(MachineReadablePrefix.INVALID.format("the body did not arrive within ${settings.bodyTimeout}")), + ) + } ?: return@admit respondLine( + call, + HttpRelayStatus.PAYLOAD_TOO_LARGE, + HttpRelayHandler.notice(MachineReadablePrefix.INVALID.format("the command exceeds $bodyCap bytes")), + ) + handler.handle(HttpRelayRequest(call.request.header(HttpHeaders.Authorization), body), Answer(call)) + } + when (verdict) { + HttpAdmission.Verdict.ADMITTED -> {} + + HttpAdmission.Verdict.CLIENT_BUSY -> { + respondLine( + call, + HttpRelayStatus.TOO_MANY_REQUESTS, + HttpRelayHandler.notice(MachineReadablePrefix.RATE_LIMITED.format("over ${settings.maxPerClient} requests at once from this client")), + ) + } + + HttpAdmission.Verdict.RELAY_BUSY -> { + respondLine( + call, + HttpRelayStatus.UNAVAILABLE, + HttpRelayHandler.notice(MachineReadablePrefix.RATE_LIMITED.format("the relay is at capacity")), + ) + } + } + } + + /** + * Who the request counts against: the peer's address, or, when the peer is a trusted proxy, the + * last address in its [HttpCommandSettings.clientAddressHeader] — the one that proxy saw. Earlier + * entries are whatever the client claimed and are never believed. + */ + private fun clientOf(call: ApplicationCall): String { + val peer = call.request.local.remoteAddress + if (peer !in settings.trustedProxies) return peer + return call.request + .header(settings.clientAddressHeader) + ?.substringAfterLast(',') + ?.trim() + ?.ifEmpty { null } ?: peer + } + + /** + * Whether `Accept-Encoding` takes gzip: listed by name or as `*`, without `q=0`. A one-line + * answer is never compressed; only a streamed one is worth it. + */ + private fun acceptsGzip(call: ApplicationCall): Boolean = + call.request + .header(HttpHeaders.AcceptEncoding) + .orEmpty() + .split(',') + .any { entry -> + val parts = entry.split(';').map { it.trim() } + val coding = parts.first().lowercase() + val q = + parts + .drop(1) + .firstOrNull { it.startsWith("q=") } + ?.removePrefix("q=") + ?.toDoubleOrNull() ?: 1.0 + (coding == "gzip" || coding == "*") && q > 0.0 + } + + /** A one-line answer: a refusal, or a command answered at once. */ + private suspend fun respondLine( + call: ApplicationCall, + status: Int, + frame: String, + ) { + when (status) { + HttpRelayStatus.UNAUTHORIZED -> call.response.header(HttpHeaders.WWWAuthenticate, WWW_AUTHENTICATE) + HttpRelayStatus.TOO_MANY_REQUESTS, HttpRelayStatus.UNAVAILABLE -> call.response.header(HttpHeaders.RetryAfter, settings.retryAfterSeconds.toString()) + } + call.respondText(frame + "\n", NDJSON, HttpStatusCode.fromValue(status)) + } + + private inner class Answer( + private val call: ApplicationCall, + ) : HttpRelayResponse { + override suspend fun single( + status: Int, + frame: String, + ) = respondLine(call, status, frame) + + /** + * Chunked: each flush puts what the handler wrote on the wire, and a full socket suspends the + * writer. Gzipped when the client takes it and the relay allows it, sync-flushed at every + * flush so the lines still arrive as they are found. + */ + override suspend fun stream(lines: suspend HttpRelayLines.() -> Unit) { + val gzip = settings.compress && acceptsGzip(call) + call.response.header(HttpHeaders.Vary, HttpHeaders.AcceptEncoding) + if (gzip) call.response.header(HttpHeaders.ContentEncoding, "gzip") + call.respondBytesWriter(NDJSON, HttpStatusCode.OK) { + val out = this + if (gzip) { + val body = GzipLines(out) + try { + object : HttpRelayLines { + override suspend fun line(frame: String) = body.line(frame) + + override suspend fun flush() = body.flush() + }.lines() + body.finish() + } finally { + body.close() + } + } else { + object : HttpRelayLines { + override suspend fun line(frame: String) { + out.writeStringUtf8(frame) + out.writeStringUtf8("\n") + } + + override suspend fun flush() = out.flush() + }.lines() + } + } + } + } + + companion object { + val NDJSON = ContentType("application", "x-ndjson") + const val REQUEST_TIMEOUT = 408 + const val WWW_AUTHENTICATE = "Nostr" + const val ACCEL_BUFFERING = "X-Accel-Buffering" + const val PREFLIGHT_MAX_AGE_SECONDS = 86_400 + } +} diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/NipFEHttpTest.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/NipFEHttpTest.kt new file mode 100644 index 0000000000..fa1f25b872 --- /dev/null +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/NipFEHttpTest.kt @@ -0,0 +1,482 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.geode + +import com.vitorpamplona.geode.server.HttpCommandSettings +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair +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.OkMessage +import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.ReqCmd +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.toHttp +import com.vitorpamplona.quartz.nip01Core.relay.server.policies.EmptyPolicy +import com.vitorpamplona.quartz.nip01Core.relay.server.policies.IRelayPolicy +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.VerifyPolicy +import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal +import com.vitorpamplona.quartz.nip01Core.store.IEventStore +import com.vitorpamplona.quartz.nip01Core.store.RawEvent +import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore +import com.vitorpamplona.quartz.nip10Notes.TextNoteEvent +import com.vitorpamplona.quartz.nip98HttpAuth.HTTPAuthorizationEvent +import com.vitorpamplona.quartz.nipFERelayOverHttp.HttpRelayClient +import com.vitorpamplona.quartz.nipFERelayOverHttp.OkHttpRelayTransport +import com.vitorpamplona.quartz.utils.TimeUtils +import kotlinx.coroutines.CompletableDeferred +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.async +import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.withTimeout +import okhttp3.MediaType.Companion.toMediaType +import okhttp3.OkHttpClient +import okhttp3.Request +import okhttp3.RequestBody.Companion.toRequestBody +import okhttp3.Response +import okio.GzipSource +import okio.buffer +import java.net.ServerSocket +import java.net.Socket +import java.util.concurrent.TimeUnit +import kotlin.test.AfterTest +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFalse +import kotlin.test.assertIs +import kotlin.test.assertTrue +import kotlin.time.Duration +import kotlin.time.Duration.Companion.milliseconds +import kotlin.time.Duration.Companion.seconds + +/** + * NIP-FE end to end: quartz's [HttpRelayClient] and raw OkHttp requests against a real [KtorRelay], + * covering the answer shape, the status table, the headers each status carries, NIP-98 sign-in + * on an AUTH-gated relay, and sharing the relay URL with NIP-86. + */ +class NipFEHttpTest { + private val http = OkHttpClient.Builder().build() + private val alice = NostrSignerInternal(KeyPair()) + private val running = mutableListOf>() + + @AfterTest + fun teardown() { + running.forEach { (server, relay) -> + server.stop(0, 1_000) + relay.close() + } + } + + /** A relay whose advertised URL is the one it listens on, so a NIP-98 `u` names a reachable endpoint. */ + private fun start( + path: String = "/", + settings: HttpCommandSettings? = HttpCommandSettings(), + policy: ((NormalizedRelayUrl) -> IRelayPolicy)? = null, + store: (IEventStore) -> IEventStore = { it }, + ): NormalizedRelayUrl { + val port = ServerSocket(0).use { it.localPort } + val url = "ws://127.0.0.1:$port$path".normalizeRelayUrl() + val events = store(EventStore(dbName = null, relay = url, indexStrategy = RelayIndexingStrategy)) + val relay = RelayEngine(url, events, policyBuilder = { policy?.invoke(url) ?: EmptyPolicy }) + val server = KtorRelay(relay, host = "127.0.0.1", port = port, path = path, httpCommands = settings).start() + running += server to relay + return url + } + + private fun client(signer: NostrSignerInternal? = null) = HttpRelayClient(OkHttpRelayTransport { http }, signer) + + private suspend fun note(text: String): Event = alice.sign(TextNoteEvent.build(text)) + + private fun post( + url: String, + body: String, + authorization: String? = null, + contentType: String = "text/plain", + ): Response = + http + .newCall( + Request + .Builder() + .url(url) + .post(body.toRequestBody(contentType.toMediaType())) + .apply { authorization?.let { header("Authorization", it) } } + .build(), + ).execute() + + private fun Response.lines() = body.string().lines().filter { it.isNotEmpty() } + + @Test + fun aReqStreamsWhatWasPublishedAndEndsOnEose() = + runBlocking { + val relay = start() + val notes = listOf(note("one"), note("two"), note("three")) + notes.forEach { assertTrue(assertIs(client().publish(relay, it).last).success) } + + val got = mutableListOf() + val answer = client().req(relay, listOf(Filter(kinds = listOf(TextNoteEvent.KIND))), onEvent = got::add) + assertEquals(200, answer.status) + assertTrue(answer.complete) + assertIs(answer.last) + assertEquals(notes.map { it.id }.toSet(), got.map { it.id }.toSet()) + } + + @Test + fun theWireIsNdjsonInTheSocketsOwnFrames() = + runBlocking { + val relay = start() + val n = note("hello") + client().publish(relay, n) + post(relay.toHttp(), """["REQ","q",{"ids":["${n.id}"]}]""").use { response -> + assertEquals(200, response.code) + assertTrue(response.header("Content-Type")!!.startsWith("application/x-ndjson")) + assertEquals("no", response.header("X-Accel-Buffering")) + assertEquals("*", response.header("Access-Control-Allow-Origin")) + val lines = response.lines() + assertEquals(listOf("""["EVENT","q",${n.toJson()}]""", """["EOSE","q"]"""), lines) + } + } + + @Test + fun aCountIsOneCountLine() = + runBlocking { + val relay = start() + client().publish(relay, note("a")) + client().publish(relay, note("b")) + val answer = client().count(relay, listOf(Filter(kinds = listOf(TextNoteEvent.KIND)))) + assertTrue(answer.complete) + assertEquals(2, assertIs(answer.last).result.count) + } + + @Test + fun aDuplicateIs200AndAForgeryIs400() = + runBlocking { + val relay = start(policy = { VerifyPolicy }) + val n = note("once") + assertEquals(200, client().publish(relay, n).status) + assertEquals(200, client().publish(relay, n).status) + + val forged = n.toJson().replace("\"once\"", "\"twice\"") + post(relay.toHttp(), """["EVENT",$forged]""").use { response -> + assertEquals(400, response.code) + val line = response.lines().single() + assertTrue(line.startsWith("""["OK","${n.id}",false,"invalid:"""), line) + } + } + + @Test + fun aBodyThatIsNotACommandIs400AndOneOverTheCapIs413() { + val relay = start(settings = HttpCommandSettings(maxBodyBytes = 64)) + for (body in listOf("hello", """{"kinds":[1]}""", """["CLOSE","q"]""")) { + post(relay.toHttp(), body).use { response -> + assertEquals(400, response.code, body) + assertTrue(response.lines().single().startsWith("""["NOTICE","invalid:""")) + } + } + post(relay.toHttp(), """["REQ","q",{"authors":["${"a".repeat(64)}"]}]""").use { response -> + assertEquals(413, response.code) + assertTrue(response.lines().single().startsWith("""["NOTICE","invalid:""")) + } + } + + @Test + fun aRateLimitedRefusalIs429WithRetryAfter() { + val relay = + start( + settings = HttpCommandSettings(retryAfterSeconds = 7), + policy = { + object : PassThroughPolicy() { + override fun accept(cmd: ReqCmd): PolicyResult = PolicyResult.Rejected("rate-limited: slow down") + } + }, + ) + post(relay.toHttp(), """["REQ","q",{}]""").use { response -> + assertEquals(429, response.code) + assertEquals("7", response.header("Retry-After")) + assertEquals("""["CLOSED","q","rate-limited: slow down"]""", response.lines().single()) + } + } + + @Test + fun aPreflightLetsAnyOriginPostWithAuthorization() { + val relay = start() + val preflight = + Request + .Builder() + .url(relay.toHttp()) + .method("OPTIONS", null) + .header("Origin", "https://app.example") + .header("Access-Control-Request-Method", "POST") + .header("Access-Control-Request-Headers", "authorization") + .build() + http.newCall(preflight).execute().use { response -> + assertEquals(204, response.code) + assertEquals("*", response.header("Access-Control-Allow-Origin")) + assertTrue(response.header("Access-Control-Allow-Methods")!!.contains("POST")) + assertTrue(response.header("Access-Control-Allow-Headers")!!.contains("Authorization")) + } + } + + @Test + fun anAuthGatedRelayAnswers401UntilANip98TokenSignsTheRequestIn() = + runBlocking { + val relay = start(policy = ::SignInPolicy) + post(relay.toHttp(), """["REQ","q",{}]""").use { response -> + assertEquals(401, response.code) + assertEquals("Nostr", response.header("WWW-Authenticate")) + assertTrue(response.lines().single().startsWith("""["CLOSED","q","auth-required:""")) + } + + val unsigned = client().req(relay, listOf(Filter(kinds = listOf(1)))) {} + assertEquals(401, unsigned.status) + assertTrue(unsigned.complete) + + val n = note("signed in") + val published = client(alice).publish(relay, n) + assertEquals(200, published.status) + assertTrue(assertIs(published.last).success) + + val got = mutableListOf() + val read = HttpRelayClient(OkHttpRelayTransport { http }, alice, signFirst = true).req(relay, listOf(Filter(ids = listOf(n.id))), onEvent = got::add) + assertEquals(200, read.status) + assertTrue(read.complete) + assertEquals(listOf(n.id), got.map { it.id }) + } + + @Test + fun aTokenForAnotherBodyDoesNotSignIn() = + runBlocking { + val relay = start(policy = ::SignInPolicy) + val url = relay.toHttp() + val signed = """["REQ","q",{"kinds":[1]}]""" + val token = alice.sign(HTTPAuthorizationEvent.build(url, "POST", signed.encodeToByteArray())).toAuthToken() + post(url, """["REQ","q",{"kinds":[0]}]""", token).use { response -> + assertEquals(401, response.code) + assertTrue(response.lines().single().contains("payload")) + } + post(url, signed, token).use { assertEquals(200, it.code) } + post(url, signed, token).use { assertEquals(200, it.code, "good again for the same body within its window") } + } + + @Test + fun aTokenSignedAtTheOnionAddressVerifies() = + runBlocking { + val onion = "ws://2gzyxa5ihm7nsggfxnu52rck2vv4rvmdlkiu3zzui5du4xyclen53wid.onion/".normalizeRelayUrl() + val relay = start(policy = ::SignInPolicy, settings = HttpCommandSettings(alternateUrls = listOf(onion))) + val body = """["REQ","q",{"kinds":[1]}]""" + val token = alice.sign(HTTPAuthorizationEvent.build(onion.toHttp(), "POST", body.encodeToByteArray())).toAuthToken() + post(relay.toHttp(), body, token).use { assertEquals(200, it.code) } + } + + @Test + fun commandsGoToTheRelayUrlPathIncluded() = + runBlocking { + val relay = start(path = "/nostr") + assertTrue(relay.toHttp().trimEnd('/').endsWith("/nostr")) + assertTrue(client().req(relay, listOf(Filter(kinds = listOf(1)))) {}.complete) + post(relay.toHttp().replace("/nostr", ""), """["REQ","q",{}]""").use { assertEquals(404, it.code) } + } + + @Test + fun nip86CallsKeepTheRelayUrlByTheirContentType() { + val relay = start() + post(relay.toHttp(), """{"method":"supportedmethods","params":[]}""", contentType = "application/nostr+json+rpc").use { response -> + assertEquals(401, response.code, "NIP-86 asks for its own NIP-98 token") + assertFalse(response.header("Content-Type")!!.startsWith("application/x-ndjson")) + } + } + + @Test + fun turnedOffEveryPostIsNip86Again() { + val relay = start(settings = null) + post(relay.toHttp(), """["REQ","q",{}]""").use { response -> + assertFalse(response.header("Content-Type")!!.startsWith("application/x-ndjson")) + assertEquals(401, response.code) + } + } + + /** Answers a filter naming [HELD_KIND] with what is stored, then holds its EOSE until [release]. */ + private fun holdingEose(release: CompletableDeferred): (IEventStore) -> IEventStore = + { real -> + object : IEventStore by real { + override suspend fun rawQuery( + filters: List, + onEach: (RawEvent) -> Unit, + ) { + real.rawQuery(filters, onEach) + if (filters.any { it.kinds?.contains(HELD_KIND) == true }) release.await() + } + } + } + + @Test + fun aGzippedAnswerStillArrivesLineByLine() = + runBlocking { + val release = CompletableDeferred() + val relay = start(store = holdingEose(release)) + val held = alice.sign(TimeUtils.now(), HELD_KIND, emptyArray(), "held") + client().publish(relay, held) + + // Asking for gzip by hand turns off OkHttp's own inflating, so the body is read as sent. + val patient = http.newBuilder().readTimeout(10, TimeUnit.SECONDS).build() + val request = + Request + .Builder() + .url(relay.toHttp()) + .post("""["REQ","q",{"kinds":[$HELD_KIND]}]""".toRequestBody("text/plain".toMediaType())) + .header("Accept-Encoding", "gzip") + .build() + patient.newCall(request).execute().use { response -> + assertEquals("gzip", response.header("Content-Encoding")) + val lines = GzipSource(response.body.source()).buffer() + // The store has not answered EOSE: this line can only be here if the gzip was flushed. + assertEquals("""["EVENT","q",${held.toJson()}]""", lines.readUtf8Line()) + assertFalse(release.isCompleted) + release.complete(Unit) + assertEquals("""["EOSE","q"]""", lines.readUtf8Line()) + assertEquals(null, lines.readUtf8Line()) + } + } + + @Test + fun theClientReadsAGzippedAnswerAsItStreams() = + runBlocking { + val release = CompletableDeferred() + val relay = start(store = holdingEose(release)) + val held = alice.sign(TimeUtils.now(), HELD_KIND, emptyArray(), "held") + client().publish(relay, held) + + val first = CompletableDeferred() + val answer = async(Dispatchers.IO) { client().req(relay, listOf(Filter(kinds = listOf(HELD_KIND))), onEvent = { first.complete(it) }) } + assertEquals(held.id, withTimeout(10.seconds) { first.await() }.id, "the event arrives before EOSE is written") + release.complete(Unit) + assertTrue(answer.await().complete) + } + + @Test + fun withoutGzipOrWithItOffTheBodyIsPlain() { + for ((settings, accept) in listOf(HttpCommandSettings() to "identity", HttpCommandSettings(compress = false) to "gzip")) { + val relay = start(settings = settings) + val request = + Request + .Builder() + .url(relay.toHttp()) + .post("""["REQ","q",{"kinds":[1]}]""".toRequestBody("text/plain".toMediaType())) + .header("Accept-Encoding", accept) + .build() + http.newCall(request).execute().use { response -> + assertEquals(null, response.header("Content-Encoding"), accept) + assertEquals(listOf("""["EOSE","q"]"""), response.lines()) + } + } + } + + /** Writes [request] on a plain socket to the relay at [relay] and returns the status line of the answer. */ + private fun rawStatus( + relay: NormalizedRelayUrl, + request: String, + ): String = + Socket( + "127.0.0.1", + relay + .toHttp() + .substringAfterLast(':') + .trimEnd('/') + .toInt(), + ).use { socket -> + socket.soTimeout = 10_000 + socket.getOutputStream().write(request.encodeToByteArray()) + socket.getOutputStream().flush() + socket.getInputStream().bufferedReader().readLine() + } + + private fun chunked(body: String) = "${body.length.toString(16)}\r\n$body\r\n0\r\n\r\n" + + @Test + fun aMalformedContentTypeIsStillACommand() { + val relay = start() + val body = """["REQ","q",{"kinds":[1]}]""" + val status = rawStatus(relay, "POST / HTTP/1.1\r\nHost: x\r\nContent-Type: garbage\r\nContent-Length: ${body.length}\r\nConnection: close\r\n\r\n$body") + assertEquals("HTTP/1.1 200 OK", status) + } + + @Test + fun aBodyThatNeverFinishesIsDroppedAndFreesItsSlot() { + val relay = start(settings = HttpCommandSettings(maxPerClient = 1, bodyTimeout = 500.milliseconds)) + val port = + relay + .toHttp() + .substringAfterLast(':') + .trimEnd('/') + .toInt() + Socket("127.0.0.1", port).use { slow -> + slow.soTimeout = 10_000 + slow.getOutputStream().write("POST / HTTP/1.1\r\nHost: x\r\nContent-Length: 100\r\n\r\n[\"REQ\"".encodeToByteArray()) + slow.getOutputStream().flush() + assertEquals("HTTP/1.1 408 Request Timeout", slow.getInputStream().bufferedReader().readLine()) + } + val body = """["REQ","q",{"kinds":[1]}]""" + assertEquals("HTTP/1.1 200 OK", rawStatus(relay, "POST / HTTP/1.1\r\nHost: x\r\nContent-Length: ${body.length}\r\nConnection: close\r\n\r\n$body")) + } + + @Test + fun aChunkedBodyGrowsItsBufferAndStopsAtTheCap() { + val relay = start(settings = HttpCommandSettings(maxBodyBytes = 20_000)) + // No Content-Length: the buffer starts small and grows past it. + val big = """["REQ","q",{"search":"${"x".repeat(10_000)}"}]""" + val head = "POST / HTTP/1.1\r\nHost: x\r\nTransfer-Encoding: chunked\r\nConnection: close\r\n\r\n" + assertEquals("HTTP/1.1 200 OK", rawStatus(relay, head + chunked(big))) + val over = """["REQ","q",{"search":"${"x".repeat(30_000)}"}]""" + assertEquals("HTTP/1.1 413 Payload Too Large", rawStatus(relay, head + chunked(over))) + } + + @Test + fun nip11AdvertisesFE() { + val relay = start() + val request = + Request + .Builder() + .url(relay.url.replace("ws://", "http://")) + .header("Accept", "application/nostr+json") + .build() + http.newCall(request).execute().use { response -> + assertTrue(response.body.string().contains("\"FE\"")) + } + } + + @Test + fun noFirstFrameWithinTheDeadlineIs503WithRetryAfter() = + runBlocking { + val relay = start(settings = HttpCommandSettings(deadline = Duration.ZERO, retryAfterSeconds = 3)) + val answer = client().req(relay, listOf(Filter(kinds = listOf(1)))) {} + assertEquals(503, answer.status) + assertEquals("3", answer.retryAfter) + assertTrue(answer.complete) + assertFalse(answer.last is EoseMessage) + } + + private companion object { + /** A regular kind (stored), not a text note, so no other test reads it. */ + const val HELD_KIND = 7_777 + } +} diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/config/StaticConfigTest.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/config/StaticConfigTest.kt index 2ba0823441..a9469e37aa 100644 --- a/geode/src/test/kotlin/com/vitorpamplona/geode/config/StaticConfigTest.kt +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/config/StaticConfigTest.kt @@ -28,6 +28,7 @@ import kotlin.test.assertEquals import kotlin.test.assertFailsWith import kotlin.test.assertNotNull import kotlin.test.assertTrue +import kotlin.time.Duration.Companion.seconds class StaticConfigTest { @Test @@ -219,7 +220,7 @@ class StaticConfigTest { assertEquals("wss://relay.example.com/", c.info.relay_url) assertEquals("Example", c.info.name) - assertEquals(listOf(1, 9, 11, 42), c.info.supported_nips) + assertEquals(listOf("1", "9", "11", "42"), c.info.supported_nips) assertEquals("127.0.0.1", c.network.host) assertEquals(9988, c.network.port) @@ -249,6 +250,56 @@ class StaticConfigTest { assertEquals(listOf("1", "11", "42"), info.document.supported_nips) } + @Test + fun supportedNipsTakeHexNamedNipsAsStrings() { + val c = + StaticConfig.fromToml( + """ + [info] + supported_nips = [1, 11, "FE"] + """.trimIndent(), + ) + assertEquals(listOf("1", "11", "FE"), c.resolveInfo().document.supported_nips) + } + + @Test + fun httpCommandsAreOnByDefaultAndAdvertised() { + val c = StaticConfig.fromToml("") + assertTrue(c.http.enabled) + assertNotNull(c.http.toSettings()) + assertTrue("FE" in c.resolveInfo().document.supported_nips!!) + } + + @Test + fun httpSectionParsesAndTurningItOffDropsFE() { + val c = + StaticConfig.fromToml( + """ + [http] + enabled = false + max_concurrent_requests = 10 + max_requests_per_client = 2 + deadline_seconds = 5 + alternate_urls = ["ws://2gzyxa5ihm7nsggfxnu52rck2vv4rvmdlkiu3zzui5du4xyclen53wid.onion/"] + trusted_proxies = ["127.0.0.1"] + """.trimIndent(), + ) + assertEquals(10, c.http.max_concurrent_requests) + assertEquals(2, c.http.max_requests_per_client) + assertEquals(listOf("127.0.0.1"), c.http.trusted_proxies) + assertEquals(null, c.http.toSettings()) + assertTrue("FE" !in c.resolveInfo().document.supported_nips!!) + val on = c.copy(http = c.http.copy(enabled = true)).http.toSettings()!! + assertEquals(5.seconds, on.deadline) + assertEquals(1, on.alternateUrls.size) + } + + @Test + fun httpLimitsMustBeSane() { + assertFailsWith { StaticConfig.fromToml("[http]\ndeadline_seconds = 0").validate() } + assertFailsWith { StaticConfig.fromToml("[http]\nmax_requests_per_client = -1").validate() } + } + @Test fun loadsTheBundledExampleConfigCleanly() { // The example file lives at the module root so operators have a diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/server/HttpAdmissionTest.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/server/HttpAdmissionTest.kt new file mode 100644 index 0000000000..ed2cbb77af --- /dev/null +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/server/HttpAdmissionTest.kt @@ -0,0 +1,71 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.geode.server + +import kotlinx.coroutines.CompletableDeferred +import kotlinx.coroutines.async +import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.yield +import kotlin.test.Test +import kotlin.test.assertEquals + +class HttpAdmissionTest { + @Test + fun aClientOverItsShareIsBusyAndARelayOverItsCapIsTooWhileOthersStillRun() = + runBlocking { + val admission = HttpAdmission(maxConcurrent = 2, maxPerClient = 1) + val release = CompletableDeferred() + val a = async { admission.admit("a") { release.await() } } + val b = async { admission.admit("b") { release.await() } } + while (admission.inFlight < 2) yield() + + assertEquals(HttpAdmission.Verdict.CLIENT_BUSY, admission.admit("a") {}) + assertEquals(HttpAdmission.Verdict.RELAY_BUSY, admission.admit("c") {}) + + release.complete(Unit) + assertEquals(HttpAdmission.Verdict.ADMITTED, a.await()) + assertEquals(HttpAdmission.Verdict.ADMITTED, b.await()) + assertEquals(0, admission.inFlight) + assertEquals(HttpAdmission.Verdict.ADMITTED, admission.admit("a") {}) + assertEquals(HttpAdmission.Verdict.ADMITTED, admission.admit("c") {}) + } + + @Test + fun zeroIsNoLimit() = + runBlocking { + val admission = HttpAdmission(maxConcurrent = 0, maxPerClient = 0) + val release = CompletableDeferred() + val held = List(50) { async { admission.admit("a") { release.await() } } } + while (admission.inFlight < 50) yield() + assertEquals(HttpAdmission.Verdict.ADMITTED, admission.admit("a") {}) + release.complete(Unit) + held.forEach { assertEquals(HttpAdmission.Verdict.ADMITTED, it.await()) } + } + + @Test + fun aFailingRequestStillLeaves() = + runBlocking { + val admission = HttpAdmission(maxConcurrent = 1, maxPerClient = 1) + runCatching { admission.admit("a") { error("boom") } } + assertEquals(0, admission.inFlight) + assertEquals(HttpAdmission.Verdict.ADMITTED, admission.admit("a") {}) + } +} 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 df77566329..adc9b47ec2 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 @@ -189,10 +189,25 @@ class RelaySession( } /** - * Dispatches an already-parsed command, for a transport that parsed and - * sized it itself (NIP-FE's HTTP bodies). [IRelayPolicy.acceptMessage] - * is not run here — it judges wire text — so such a caller applies the - * message-length limit before calling. + * [receive] for a transport that already parsed [parsed] out of [command] + * (NIP-FE's HTTP bodies): [IRelayPolicy.acceptMessage] still judges the + * wire text, as it would on the socket, but the frame is not parsed twice. + */ + suspend fun receive( + command: String, + parsed: Command, + ) { + policy.acceptMessage(command)?.let { reason -> + send(NoticeMessage(reason)) + return + } + receive(parsed) + } + + /** + * Dispatches an already-parsed command. [IRelayPolicy.acceptMessage] is + * not run here — it judges wire text — so a caller that has the text + * uses the overload that takes both. */ suspend fun receive(cmd: Command) { if (!cmd.isValid()) { diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayAnswer.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayAnswer.kt new file mode 100644 index 0000000000..98b04bfd88 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayAnswer.kt @@ -0,0 +1,84 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.quartz.nipFERelayOverHttp + +import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message + +/** + * A NIP-FE answer as the client got it. [complete] is false when the body stopped before the frame + * that ends the command's answer: what came is a prefix, and a REQ's prefix looks exactly like a + * short result, so a caller MUST NOT read an incomplete answer as the whole one. + */ +class HttpRelayAnswer( + val status: Int, + /** The frame the answer ended on (EOSE, CLOSED, COUNT, OK, NOTICE), or the last one read when cut off. */ + val last: Message?, + /** Whether the answer reached its end, rather than being cut off. */ + val complete: Boolean, + /** The `Retry-After` the relay sent with a 429 or 503. */ + val retryAfter: String? = null, +) + +/** + * NIP-FE, client side: reads an answer one line at a time as it arrives and keeps track of whether + * it ended where [command]'s answer ends. A 200 streams up to that frame; any other status is one + * refusal line, which is the whole answer whatever frame it is. + */ +class HttpRelayAnswerReader( + val command: HttpRelayCommand, + val status: Int, +) { + var last: Message? = null + private set + + private var lines = 0 + private var ended = false + private var broken = false + + /** + * The frame on [line], or null for a blank line. A line that does not parse, or anything after + * the answer ended, marks the answer incomplete: the body is not what this NIP says it is. + */ + fun read(line: String): Message? { + if (line.isBlank()) return null + if (ended || broken) { + broken = true + return null + } + val message = + try { + OptimizedJsonMapper.fromJsonToMessage(line) + } catch (_: Exception) { + broken = true + return null + } + lines++ + last = message + ended = status != HttpRelayStatus.OK || command.ends(message) + return message + } + + /** Whether the body read so far is a whole answer. Asked once the body ends. */ + val complete: Boolean get() = ended && !broken && (status == HttpRelayStatus.OK || lines == 1) + + fun answer(retryAfter: String? = null) = HttpRelayAnswer(status, last, complete, retryAfter) +} diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayClient.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayClient.kt new file mode 100644 index 0000000000..af7eb13305 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayClient.kt @@ -0,0 +1,127 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.quartz.nipFERelayOverHttp + +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.relay.client.single.newSubId +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.EventMessage +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message +import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.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.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.toHttp +import com.vitorpamplona.quartz.nip01Core.signers.NostrSigner +import com.vitorpamplona.quartz.nip98HttpAuth.HTTPAuthorizationEvent + +/** + * NIP-FE client: one relay command per request, POSTed to the relay's URL as the frame the + * websocket would carry, its answer read line by line as the relay writes it, in the socket's own + * frames. Nothing stays open after a call returns. [transport] carries the bytes (OkHttp on + * JVM/Android: `OkHttpRelayTransport`), as a websocket builder does for the NostrClient. + * + * With a [signer], a request the relay refuses with 401 goes once more carrying a NIP-98 token for + * its exact body, as a websocket client answers a NIP-42 challenge; without one, the 401 is the + * answer. [signFirst] sends the token with the first try instead, saving the round trip against a + * relay known to want it. + */ +class HttpRelayClient( + private val transport: HttpRelayTransport, + private val signer: NostrSigner? = null, + private val signFirst: Boolean = false, +) { + /** Stored events matching [filters], each to [onEvent] as it arrives. Complete only if it ended on EOSE or CLOSED. */ + suspend fun req( + relay: NormalizedRelayUrl, + filters: List, + subId: String = newSubId(), + onEvent: (Event) -> Unit, + ): HttpRelayAnswer = + send(relay, ReqCmd(subId, filters)) { + if (it is EventMessage) onEvent(it.event) + } + + /** NIP-45: [HttpRelayAnswer.last] is the COUNT, or the refusal. */ + suspend fun count( + relay: NormalizedRelayUrl, + filters: List, + queryId: String = newSubId(), + ): HttpRelayAnswer = send(relay, CountCmd(queryId, filters)) + + /** Publishes [event]: [HttpRelayAnswer.last] is its OK, or the refusal. */ + suspend fun publish( + relay: NormalizedRelayUrl, + event: Event, + ): HttpRelayAnswer = send(relay, EventCmd(event)) + + /** Posts [cmd] (a REQ, COUNT or EVENT) to [relay], handing every frame to [onMessage] as it is read. */ + suspend fun send( + relay: NormalizedRelayUrl, + cmd: Command, + onMessage: (Message) -> Unit = {}, + ): HttpRelayAnswer { + val command = requireNotNull(HttpRelayCommand.of(cmd)) { "NIP-FE carries REQ, COUNT and EVENT, not ${cmd.label()}" } + val url = relay.toHttp() + val body = cmd.toJson().encodeToByteArray() + if (signer == null || signFirst) return post(relay, command, url, body, token(url, body), onMessage, retrying = false) + val first = post(relay, command, url, body, null, onMessage, retrying = true) + if (first.status != HttpRelayStatus.UNAUTHORIZED) return first + return post(relay, command, url, body, token(url, body), onMessage, retrying = false) + } + + private suspend fun token( + url: String, + body: ByteArray, + ): String? = signer?.sign(HTTPAuthorizationEvent.build(url, "POST", body))?.toAuthToken() + + private suspend fun post( + relay: NormalizedRelayUrl, + command: HttpRelayCommand, + url: String, + body: ByteArray, + authorization: String?, + onMessage: (Message) -> Unit, + /** A signed try follows a 401, so that refusal is not the answer and is not handed on. */ + retrying: Boolean, + ): HttpRelayAnswer { + var reader: HttpRelayAnswerReader? = null + var retryAfter: String? = null + var deliver = true + transport.post( + relay = relay, + url = url, + body = body, + authorization = authorization, + onStatus = { status, after -> + reader = HttpRelayAnswerReader(command, status) + retryAfter = after + deliver = !(retrying && status == HttpRelayStatus.UNAUTHORIZED) + }, + onLine = { line -> + val message = checkNotNull(reader) { "a line before the status" }.read(line) + if (message != null && deliver) onMessage(message) + }, + ) + return checkNotNull(reader) { "the transport returned without a status" }.answer(retryAfter) + } +} 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 cb5a865ea6..35b074fcb8 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayCommand.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayCommand.kt @@ -26,50 +26,22 @@ 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 /** - * NIP-FE: the client commands HTTP carries, one path each. A body is the command's arguments - * after its subscription id (a lone object where the command takes one); the answer ends on the - * first frame [ends] accepts. + * NIP-FE: the client commands an HTTP request may carry, each as the frame a client sends on the + * websocket, and the frame that ends each one's answer. */ -enum class HttpRelayCommand( - val path: String, -) { - REQ("/req"), - COUNT("/count"), - EVENT("/event"), +enum class HttpRelayCommand { + REQ, + COUNT, + EVENT, ; - /** - * The client frame [body] stands for, or null when it plainly is not this command's arguments. - * The body is spliced in as sent and the engine parses the frame, as it parses socket text, so - * any other malformed body is the engine's NOTICE. The verb and subscription id come first and - * the parser reads one value, so nothing a body holds can make it another command. - */ - fun frameOf(body: String): String? { - val 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]" - } - - EVENT -> { - if (!text.startsWith('{')) return null - "[\"${EventCmd.LABEL}\",$text]" - } - } - } - - /** Whether [message] is the last frame of this command's answer. */ + /** Whether [message] is the last frame of this command's answer. A NOTICE ends any: the command never ran. */ fun ends(message: Message): Boolean = message is NoticeMessage || when (this) { @@ -79,12 +51,25 @@ enum class HttpRelayCommand( } companion object { - /** - * 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" + /** The kind of [cmd], or null for one HTTP does not carry (AUTH, CLOSE, NEG-*). */ + fun of(cmd: Command): HttpRelayCommand? = + when (cmd) { + is ReqCmd -> REQ + is CountCmd -> COUNT + is EventCmd -> EVENT + else -> null + } - fun forPath(path: String): HttpRelayCommand? = entries.firstOrNull { it.path == path } + /** The frame that refuses [cmd] with [reason], as the socket would: CLOSED for a REQ or COUNT, OK false for an EVENT. */ + fun refusal( + cmd: Command, + reason: String, + ): Message = + when (cmd) { + is EventCmd -> OkMessage(cmd.event.id, false, reason) + is ReqCmd -> ClosedMessage(cmd.subId, reason) + is CountCmd -> ClosedMessage(cmd.queryId, reason) + else -> NoticeMessage(reason) + } } } 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 a4054bbd08..0b845206bf 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandler.kt @@ -21,13 +21,16 @@ package com.vitorpamplona.quartz.nipFERelayOverHttp import com.vitorpamplona.quartz.nip01Core.core.HexKey +import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.AuthMessage -import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.ClosedMessage import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.MachineReadablePrefix import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.NoticeMessage +import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.Command import com.vitorpamplona.quartz.nip01Core.relay.server.RelayServerBase import com.vitorpamplona.quartz.nip01Core.relay.server.SessionSink import com.vitorpamplona.quartz.nip98HttpAuth.Nip98AuthVerifier +import com.vitorpamplona.quartz.nip98HttpAuth.tags.UrlTag import kotlinx.coroutines.CancellationException import kotlinx.coroutines.CompletableDeferred import kotlinx.coroutines.TimeoutCancellationException @@ -36,18 +39,19 @@ 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]. */ +/** One NIP-FE request as the handler needs it: a POST to the relay's URL that is not a NIP-86 call. */ 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. + * The body: one client frame. 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, ) @@ -82,16 +86,17 @@ interface HttpRelayLines { 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, without their subscription id. Nothing outlives the request. + * NIP-FE: one relay command per HTTP request. The body is the frame a client would send on the + * websocket; it runs on its own [RelayServerBase] session, fed as socket text, so every limit and + * policy the socket applies applies here; and the answer is the session's frames as the socket + * would carry them, up to the one that ends 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 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. */ + /** The URLs a NIP-98 `u` may name (the relay's http URL, its .onion), asked per request; never from the request. */ private val origins: () -> List, /** How long one answer may run, first byte to last. [Duration.INFINITE] turns the deadline off. */ private val deadline: Duration = DEFAULT_DEADLINE, @@ -107,36 +112,35 @@ class HttpRelayHandler( request: HttpRelayRequest, response: HttpRelayResponse, ) { - val command = request.command val max = server.limits?.maxMessageLength + val tooLarge = "invalid: the command exceeds $max characters" maxBodyBytes?.let { cap -> - if (request.body.size > cap) { - return response.single(HttpRelayStatus.PAYLOAD_TOO_LARGE, closed("invalid: the command exceeds $max characters")) - } + if (request.body.size > cap) return response.single(HttpRelayStatus.PAYLOAD_TOO_LARGE, notice(tooLarge)) } - val frame = - command.frameOf(request.body.decodeToString()) - ?: return response.single(HttpRelayStatus.BAD_REQUEST, closed("invalid: the body is not ${command.name}'s arguments")) + val text = request.body.decodeToString() // 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")) - } + if (max != null && text.length > max) return response.single(HttpRelayStatus.PAYLOAD_TOO_LARGE, notice(tooLarge)) + + // Parsed here, once: to refuse the commands HTTP does not carry, to know what ends the + // answer and how to refuse it, and for the session, which still runs the policies that + // judge the raw text before it dispatches the parsed command. + val cmd = + try { + OptimizedJsonMapper.fromJsonToCommand(text) + } catch (_: Exception) { + null + } + val command = + cmd?.let { HttpRelayCommand.of(it) } + ?: return response.single(HttpRelayStatus.BAD_REQUEST, notice("invalid: the body is not one REQ, COUNT or EVENT frame")) + val signedIn = when (val proof = proofOf(request)) { - is Proof.Anonymous -> { - null - } - - is Proof.Signed -> { - proof.pubkey - } - - is Proof.Refused -> { - val reason = proof.reason - return response.single(HttpRelayStatus.forReason(reason), closed(reason)) - } + is Proof.Anonymous -> null + is Proof.Signed -> proof.pubkey + is Proof.Refused -> return response.single(HttpRelayStatus.forReason(proof.reason), HttpRelayCommand.refusal(cmd, proof.reason).toJson()) } - exchange(frame, command, signedIn, response) + exchange(text, cmd, command, signedIn, response) } /** A frame as queued: its wire text, its type when the engine built one, and whether it ends the answer. */ @@ -147,7 +151,8 @@ class HttpRelayHandler( ) private suspend fun exchange( - frame: String, + text: String, + cmd: Command, command: HttpRelayCommand, signedIn: HexKey?, response: HttpRelayResponse, @@ -165,23 +170,25 @@ class HttpRelayHandler( } } - fun fail(reason: String) = offer(Frame(closed(reason), ClosedMessage(HttpRelayCommand.SUB_ID, reason), last = true)) + fun refusal(reason: String) = HttpRelayCommand.refusal(cmd, reason) + + fun fail(reason: String) = refusal(reason).let { offer(Frame(it.toJson(), it, last = true)) } val sink = object : SessionSink { override fun message(message: Message) { // The challenge every connection opens with; this one proves its key by NIP-98 instead. if (message is AuthMessage) return - offer(Frame(withoutSubId(message.toJson()), message, command.ends(message))) + offer(Frame(message.toJson(), message, command.ends(message))) } - override fun raw(json: String) = offer(Frame(withoutSubId(json), null, last = false)) + override fun raw(json: String) = offer(Frame(json, null, last = false)) } val session = launch { try { server.serve(sink) { session -> val refused = signedIn?.let { session.authenticateByTransport(it) } - if (refused != null) fail(refused) else session.receive(frame) + if (refused != null) fail(refused) else session.receive(text, cmd) ended.await() } } catch (e: CancellationException) { @@ -202,7 +209,7 @@ class HttpRelayHandler( val status = HttpRelayStatus.of(first?.message) when { first == null -> { - single(HttpRelayStatus.UNAVAILABLE, closed("error: no answer within $deadline")) + single(HttpRelayStatus.UNAVAILABLE, notice("error: no answer within $deadline")) } first.last || status != HttpRelayStatus.OK -> { @@ -214,8 +221,8 @@ class HttpRelayHandler( 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")) + Ending.DEADLINE -> line(refusal("error: the answer ran past $deadline").toJson()) + Ending.CUT -> line(refusal("error: slow reader, over $maxQueuedFrames frames waiting").toJson()) } flush() } @@ -284,56 +291,55 @@ class HttpRelayHandler( } /** - * 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. + * A NIP-98 header. Its `u` may name any address in [origins], with or without the trailing + * slash, so a token signed at the .onion verifies there; the token is checked once, against the + * address it names. It must bind the body's hash and be within 60 seconds of now; within that + * window it may come again for the same body. Another scheme (a proxy's Basic, a client's + * Bearer) is not addressed to the relay and is ignored. */ private suspend fun proofOf(request: HttpRelayRequest): Proof { val header = request.authorization?.trim().orEmpty() val scheme = Nip98AuthVerifier.SCHEME if (!header.regionMatches(0, scheme, 0, scheme.length, ignoreCase = true)) return Proof.Anonymous val token = scheme + header.substring(scheme.length).trim() - val accepted = origins().map { it.trimEnd('/') + request.command.path } - if (accepted.isEmpty()) return Proof.Refused(MachineReadablePrefix.AUTH_REQUIRED.format("this relay names no url to sign")) + val addresses = origins() + if (addresses.isEmpty()) return Proof.Refused(MachineReadablePrefix.AUTH_REQUIRED.format("this relay names no url to sign")) + // The address the token names, when it is one of ours; otherwise the first, and the + // verifier refuses the mismatch (or whatever else is wrong with the token) itself. + val signed = claimedUrl(token) + val url = signed?.takeIf { u -> addresses.any { it.trimEnd('/') == u.trimEnd('/') } } ?: addresses.first() // 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 - } + // single-use. Every command it can sign is idempotent, so a repeat only repeats a read or + // re-sends an event the relay has, and a client may retry without signing again. + return when (val r = Nip98AuthVerifier(toleranceSeconds = TOKEN_WINDOW_SECONDS).verify(token, "POST", url, request.body)) { + is Nip98AuthVerifier.Result.Verified -> Proof.Signed(r.pubkey) + is Nip98AuthVerifier.Result.Missing -> Proof.Anonymous + is Nip98AuthVerifier.Result.Malformed -> Proof.Refused(MachineReadablePrefix.AUTH_REQUIRED.format("NIP-98 ${r.reason}")) } - return Proof.Refused(MachineReadablePrefix.AUTH_REQUIRED.format("NIP-98 ${refusal?.reason}")) } - private fun closed(reason: String) = withoutSubId(ClosedMessage(HttpRelayCommand.SUB_ID, reason).toJson()) + /** The `u` a NIP-98 [token] names, or null when it does not decode; the verifier then says why. */ + @OptIn(ExperimentalEncodingApi::class) + private fun claimedUrl(token: String): String? = + try { + val event = OptimizedJsonMapper.fromJson(Base64.decode(token.substring(Nip98AuthVerifier.SCHEME.length)).decodeToString()) + event.tags.firstNotNullOfOrNull(UrlTag::parse) + } catch (_: Exception) { + null + } companion object { + /** A `NOTICE` line: how a request is refused before its command runs (400, 413, 429, 503), by the handler or its host. */ + fun notice(reason: String) = NoticeMessage(reason).toJson() + val DEFAULT_DEADLINE = 30_000.milliseconds /** The websocket's slow-consumer bound in the reference relays. */ const val DEFAULT_MAX_QUEUED_FRAMES = 8192 val DEFAULT_TAIL_GRACE = 5_000.milliseconds + + /** NIP-FE: a token is good for 60 seconds either side of its `created_at`. */ + const val TOKEN_WINDOW_SECONDS = 60L } } - -/** 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/HttpRelayTransport.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayTransport.kt new file mode 100644 index 0000000000..918bb8e921 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayTransport.kt @@ -0,0 +1,47 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.quartz.nipFERelayOverHttp + +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl + +/** + * How [HttpRelayClient] reaches a relay over HTTP: one POST, its response read line by line. The + * NIP-FE side of the exchange (which URL, what body, signing, reading the answer) stays in the + * client; an implementation only moves bytes, the way a + * [com.vitorpamplona.quartz.nip01Core.relay.sockets.WebsocketBuilder] does for [com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient]. + */ +interface HttpRelayTransport { + /** + * POSTs [body] to [url], [relay]'s HTTP URL, with an `Authorization` header when [authorization] + * is not null. Calls [onStatus] once with the status and any `Retry-After`, then [onLine] for each + * body line as it arrives, and returns when the body ends. A connection that drops mid-body + * returns normally: the answer reader tells a cut-off answer from a whole one. Cancelling the + * caller cancels the request. + */ + suspend fun post( + relay: NormalizedRelayUrl, + url: String, + body: ByteArray, + authorization: String?, + onStatus: (status: Int, retryAfter: String?) -> Unit, + onLine: (String) -> Unit, + ) +} diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayAnswerReaderTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayAnswerReaderTest.kt new file mode 100644 index 0000000000..085290682e --- /dev/null +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayAnswerReaderTest.kt @@ -0,0 +1,101 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.quartz.nipFERelayOverHttp + +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.ClosedMessage +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.CountMessage +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.EoseMessage +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.EventMessage +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.OkMessage +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFalse +import kotlin.test.assertIs +import kotlin.test.assertNull +import kotlin.test.assertTrue + +/** NIP-FE, client side: an answer's lines back into frames, and a whole answer told from a cut one. */ +class HttpRelayAnswerReaderTest { + private val event = + """{"id":"a8e0b2f2c1e7e4f5b0a1f9c1b3d2e6a7c8b9d0e1f2a3b4c5d6e7f8091a2b3c4d","pubkey":"79be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798",""" + + """"created_at":1700000000,"kind":1,"tags":[],"content":"hi","sig":"${"0".repeat(128)}"}""" + + private fun read( + command: HttpRelayCommand, + status: Int, + vararg lines: String, + ) = HttpRelayAnswerReader(command, status).also { reader -> lines.forEach { reader.read(it) } } + + @Test + fun aReqThatEndsOnEoseIsComplete() { + val reader = read(HttpRelayCommand.REQ, 200, """["EVENT","q",$event]""", """["EOSE","q"]""") + assertTrue(reader.complete) + assertIs(reader.last) + } + + @Test + fun aReqWithItsTailMissingIsIncomplete() { + val reader = HttpRelayAnswerReader(HttpRelayCommand.REQ, 200) + val message = reader.read("""["EVENT","q",$event]""") + assertIs(message) + assertEquals("hi", message.event.content) + assertEquals("q", message.subId) + assertFalse(reader.complete) + assertFalse(HttpRelayAnswerReader(HttpRelayCommand.REQ, 200).complete, "an empty body is no answer") + } + + @Test + fun aClosedEndsAReqAndACountButNotAnEvent() { + assertTrue(read(HttpRelayCommand.REQ, 200, """["EVENT","q",$event]""", """["CLOSED","q","error: the answer ran past 30s"]""").complete) + assertTrue(read(HttpRelayCommand.COUNT, 200, """["CLOSED","q","error: x"]""").complete) + assertFalse(read(HttpRelayCommand.EVENT, 200, """["CLOSED","q","error: x"]""").complete) + } + + @Test + fun aCountAndAnOkAreWholeAnswers() { + val count = read(HttpRelayCommand.COUNT, 200, """["COUNT","c",{"count":7}]""") + assertTrue(count.complete) + assertEquals(7, assertIs(count.last).result.count) + val ok = read(HttpRelayCommand.EVENT, 200, """["OK","abc",true,""]""") + assertTrue(ok.complete) + assertTrue(assertIs(ok.last).success) + } + + @Test + fun aRefusalIsOneLineWhateverItsFrame() { + val refused = read(HttpRelayCommand.EVENT, 401, """["CLOSED","q","auth-required: sign in"]""") + assertTrue(refused.complete) + assertEquals("auth-required: sign in", assertIs(refused.last).message) + assertFalse(read(HttpRelayCommand.REQ, 403, """["CLOSED","q","blocked: no"]""", """["EOSE","q"]""").complete) + } + + @Test + fun anythingAfterTheEndOrALineThatIsNoFrameBreaksTheAnswer() { + assertFalse(read(HttpRelayCommand.REQ, 200, """["EOSE","q"]""", """["EVENT","q",$event]""").complete) + assertFalse(read(HttpRelayCommand.REQ, 200, """["EVENT","q",$event]""", """["EOSE","q""").complete) + assertTrue(read(HttpRelayCommand.REQ, 200, """["EOSE","q"]""", "", " ").complete, "blank lines are not frames") + } + + @Test + fun blankLinesAreSkipped() { + assertNull(HttpRelayAnswerReader(HttpRelayCommand.REQ, 200).read("")) + } +} diff --git a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/OkHttpRelayTransport.kt b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/OkHttpRelayTransport.kt new file mode 100644 index 0000000000..ee65bda983 --- /dev/null +++ b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/OkHttpRelayTransport.kt @@ -0,0 +1,92 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.quartz.nipFERelayOverHttp + +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.awaitCancellation +import kotlinx.coroutines.coroutineScope +import kotlinx.coroutines.launch +import kotlinx.coroutines.withContext +import okhttp3.MediaType.Companion.toMediaType +import okhttp3.OkHttpClient +import okhttp3.Request +import okhttp3.RequestBody.Companion.toRequestBody +import okhttp3.coroutines.executeAsync +import okio.IOException + +/** + * [HttpRelayTransport] over OkHttp. [httpClient] picks the client per relay, as + * [com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.BasicOkHttpWebSocket.Builder] does for + * websockets, so a .onion relay can go through Tor and a clearnet one directly. OkHttp asks for + * gzip and inflates it as the body streams, so a relay's sync-flushed gzip arrives line by line. + */ +class OkHttpRelayTransport( + private val httpClient: (NormalizedRelayUrl) -> OkHttpClient, +) : HttpRelayTransport { + override suspend fun post( + relay: NormalizedRelayUrl, + url: String, + body: ByteArray, + authorization: String?, + onStatus: (status: Int, retryAfter: String?) -> Unit, + onLine: (String) -> Unit, + ) = coroutineScope { + val request = + Request + .Builder() + .url(url) + .post(body.toRequestBody(JSON)) + .header("Accept", NDJSON) + .apply { authorization?.let { header("Authorization", it) } } + .build() + val call = httpClient(relay).newCall(request) + // A blocking read does not see coroutine cancellation; cancelling the call unblocks it. + val watcher = + launch { + try { + awaitCancellation() + } finally { + call.cancel() + } + } + try { + call.executeAsync().use { response -> + onStatus(response.code, response.header("Retry-After")) + withContext(Dispatchers.IO) { + val source = response.body.source() + try { + while (true) onLine(source.readUtf8Line() ?: break) + } catch (_: IOException) { + // The connection dropped mid-answer: what came is what the reader says it is. + } + } + } + } finally { + watcher.cancel() + } + } + + companion object { + const val NDJSON = "application/x-ndjson" + private val JSON = "application/json".toMediaType() + } +} 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 31695d8da3..5612fd3c81 100644 --- a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayHandlerTest.kt @@ -160,10 +160,15 @@ class HttpRelayHandlerTest { ) = HttpRelayHandler(MemoryRelay(backend, signedInOnly), origins = { listOf(origin) }, deadline = deadline) private fun HttpRelayHandler.ask( - command: HttpRelayCommand, - body: String, + frame: String, authorization: String? = null, - ) = Recorded().also { runBlocking { handle(HttpRelayRequest(command, authorization, body.encodeToByteArray()), it) } } + ) = Recorded().also { runBlocking { handle(HttpRelayRequest(authorization, frame.encodeToByteArray()), it) } } + + private fun req(filter: String) = """["REQ","q",$filter]""" + + private fun count(filter: String) = """["COUNT","c",$filter]""" + + private fun publish(event: Event) = """["EVENT",${event.toJson()}]""" private fun note( content: String, @@ -171,19 +176,19 @@ class HttpRelayHandlerTest { ) = alice.sign(at, 1, emptyArray(), content) private fun token( - command: HttpRelayCommand, - body: String, - ) = alice.sign(HTTPAuthorizationEvent.build(origin + command.path, "POST", body.encodeToByteArray(), System.currentTimeMillis() / 1000) {}).toAuthToken() + frame: String, + url: String = origin, + ) = alice.sign(HTTPAuthorizationEvent.build(url, "POST", frame.encodeToByteArray(), System.currentTimeMillis() / 1000) {}).toAuthToken() @Test fun aReqStreamsItsEventsAndEndsOnEose() { val a = note("a") val b = note("b") backend.events += listOf(a, b) - val answer = handler().ask(HttpRelayCommand.REQ, """{"kinds":[1]}""") + val answer = handler().ask(req("""{"kinds":[1]}""")) assertEquals(200, answer.status) assertTrue(answer.streamed) - assertEquals("""["EOSE"]""", answer.lines.last()) + assertEquals("""["EOSE","q"]""", answer.lines.last()) assertEquals( setOf(a.id, b.id), answer.lines @@ -193,23 +198,31 @@ class HttpRelayHandlerTest { ) } + @Test + fun framesGoOutAsTheSocketSendsThemWithTheClientsSubscriptionId() { + val found = note("found") + backend.events += found + val answer = handler().ask("""["REQ","mine",{"kinds":[1]}]""") + assertEquals(listOf("""["EVENT","mine",${found.toJson()}]""", """["EOSE","mine"]"""), answer.lines) + } + @Test fun anEmptyReqIsOneEoseLine() { - val answer = handler().ask(HttpRelayCommand.REQ, """[{"kinds":[30000]}]""") + val answer = handler().ask(req("""{"kinds":[30000]}""")) assertEquals(200, answer.status) - assertEquals(listOf("""["EOSE"]"""), answer.lines) + assertEquals(listOf("""["EOSE","q"]"""), answer.lines) } @Test fun anEventIsAnsweredByItsOkAndAForgeryIsRefused() { val posted = note("posted") - val ok = handler().ask(HttpRelayCommand.EVENT, posted.toJson()) + val ok = handler().ask(publish(posted)) assertEquals(200, ok.status) assertEquals(listOf("""["OK","${posted.id}",true,""]"""), ok.lines) assertTrue(backend.events.any { it.id == posted.id }) val forged = Event(posted.id, posted.pubKey, posted.createdAt, posted.kind, posted.tags, "tampered", posted.sig) - val refused = handler().ask(HttpRelayCommand.EVENT, forged.toJson()) + val refused = handler().ask(publish(forged)) assertEquals(400, refused.status, refused.lines.toString()) assertTrue(refused.lines.single().startsWith("""["OK","${forged.id}",false,"""), refused.lines.toString()) } @@ -217,76 +230,89 @@ class HttpRelayHandlerTest { @Test fun aCountIsOneCountLine() { backend.events += listOf(note("a"), note("b"), note("c")) - val answer = handler().ask(HttpRelayCommand.COUNT, """{"kinds":[1]}""") + val answer = handler().ask(count("""{"kinds":[1]}""")) assertEquals(200, answer.status) - assertTrue(answer.lines.single().startsWith("""["COUNT",{"count":3"""), answer.lines.toString()) + assertTrue(answer.lines.single().startsWith("""["COUNT","c",{"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.COUNT to "not json", - )) { - val answer = handler().ask(command, body) - assertEquals(400, answer.status, "$command '$body'") - assertTrue(answer.lines.single().startsWith("""["CLOSED","invalid:"""), answer.lines.toString()) + fun aBodyThatIsNotAReqCountOrEventIsA400AndOneOverTheLimitA413() { + for (body in listOf("[]", """{"kinds":[1]}""", "not json", """["CLOSE","q"]""", """["AUTH",{"id":"x"}]""", """["NEG-CLOSE","n"]""")) { + val answer = handler().ask(body) + assertEquals(400, answer.status, body) + assertTrue(answer.lines.single().startsWith("""["NOTICE","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) + // A REQ the engine refuses as a command: its own NOTICE, the command never ran. + val empty = handler().ask("""["REQ","",{"kinds":[1]}]""") + assertEquals(400, empty.status) + assertTrue(empty.lines.single().startsWith("""["NOTICE","""), empty.lines.toString()) + // Over the relay's message length, in characters, as the socket measures it. + val big = handler().ask(req("""{"search":"${"x".repeat(4_096)}"}""")) + assertEquals(413, big.status) + assertTrue(big.lines.single().startsWith("""["NOTICE","invalid:"""), big.lines.toString()) } @Test fun aNip98SignatureSignsTheSessionInAndAnotherSchemeDoesNot() { backend.events += note("gated") val gated = handler(signedInOnly = true) - val body = """{"kinds":[1]}""" + val frame = req("""{"kinds":[1]}""") - val anonymous = gated.ask(HttpRelayCommand.REQ, body) + val anonymous = gated.ask(frame) assertEquals(401, anonymous.status) - assertTrue(anonymous.lines.single().startsWith("""["CLOSED","auth-required:""")) + assertTrue(anonymous.lines.single().startsWith("""["CLOSED","q","auth-required:""")) - assertEquals(401, gated.ask(HttpRelayCommand.REQ, body, "Basic dXNlcjpwYXNz").status, "Basic is not addressed to the relay") + assertEquals(401, gated.ask(frame, "Basic dXNlcjpwYXNz").status, "Basic is not addressed to the relay") - val signed = gated.ask(HttpRelayCommand.REQ, body, token(HttpRelayCommand.REQ, body)) + val signed = gated.ask(frame, token(frame)) assertEquals(200, signed.status, signed.lines.toString()) - assertEquals("""["EOSE"]""", signed.lines.last()) + assertEquals("""["EOSE","q"]""", signed.lines.last()) } @Test - fun aTokenSignsOnlyItsBodyAndIsNotSingleUse() { + fun aTokenSignsOnlyItsBodyAndIsGoodAgainWithinItsWindow() { val h = handler() - val body = """{"kinds":[1]}""" - 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)) + val frame = req("""{"kinds":[1]}""") + val signed = token(frame) + assertEquals(200, h.ask(frame, signed).status) + assertEquals(200, h.ask(frame, signed).status, "the same body again only repeats the read") + val other = h.ask(req("""{"kinds":[0]}"""), signed) assertEquals(401, other.status) assertTrue("payload" in other.lines.single(), other.lines.toString()) + assertTrue(other.lines.single().startsWith("""["CLOSED","q","auth-required:"""), other.lines.toString()) + } + + @Test + fun aTokenOutsideItsSixtySecondsIsRefused() { + val frame = req("""{"kinds":[1]}""") + val stale = alice.sign(HTTPAuthorizationEvent.build(origin, "POST", frame.encodeToByteArray(), System.currentTimeMillis() / 1000 - 120) {}).toAuthToken() + assertEquals(401, handler().ask(frame, stale).status) + } + + @Test + fun anEventRefusedForItsTokenIsAnOkFalse() { + val posted = publish(note("unsigned")) + val answer = handler().ask(posted, token(req("{}"))) + assertEquals(401, answer.status) + assertTrue(answer.lines.single().startsWith("""["OK","""), answer.lines.toString()) + assertTrue(""",false,"auth-required:""" in answer.lines.single(), answer.lines.toString()) } @Test fun noFirstFrameWithinTheDeadlineIsA503() { - val answer = handler(deadline = 300.milliseconds).ask(HttpRelayCommand.REQ, """{"kinds":[$STALLED_KIND]}""") + val answer = handler(deadline = 300.milliseconds).ask(req("""{"kinds":[$STALLED_KIND]}""")) assertEquals(503, answer.status) - assertTrue(answer.lines.single().startsWith("""["CLOSED","error: no answer""")) + assertTrue(answer.lines.single().startsWith("""["NOTICE","error: no answer""")) } @Test fun aDeadlineMidAnswerEndsOnAClosedLine() { val found = note("found") backend.events += found - val answer = handler(deadline = 300.milliseconds).ask(HttpRelayCommand.REQ, """{"kinds":[1,$TRICKLE_KIND]}""") + val answer = handler(deadline = 300.milliseconds).ask(req("""{"kinds":[1,$TRICKLE_KIND]}""")) assertEquals(200, answer.status) assertTrue(found.id in answer.lines.first()) - assertTrue(answer.lines.last().startsWith("""["CLOSED","error: the answer ran past"""), answer.lines.toString()) + assertTrue(answer.lines.last().startsWith("""["CLOSED","q","error: the answer ran past"""), answer.lines.toString()) } @Test @@ -308,7 +334,7 @@ class HttpRelayHandlerTest { }.lines() } assertFailsWith { - runBlocking { h.handle(HttpRelayRequest(HttpRelayCommand.REQ, null, """{"kinds":[1,$TRICKLE_KIND]}""".encodeToByteArray()), stalled) } + runBlocking { h.handle(HttpRelayRequest(null, req("""{"kinds":[1,$TRICKLE_KIND]}""").encodeToByteArray()), stalled) } } } @@ -322,26 +348,17 @@ 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) + for (body in listOf(req("""{"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(body) assertEquals(400, answer.status, body.take(20)) } } @Test fun aFullAuthPolicyRefusesTransportSignInUntilItOptsIn() { - val body = """{"kinds":[1]}""" + val frame = req("""{"kinds":[1]}""") val relayUrl = RelayUrlNormalizer.normalize("wss://relay.example") val refusing = MemoryRelay(backend, { @@ -349,9 +366,9 @@ class HttpRelayHandlerTest { 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)) + val refused = HttpRelayHandler(refusing, origins = { listOf(origin) }).ask(frame, token(frame)) assertEquals(403, refused.status, refused.lines.toString()) - assertTrue(refused.lines.single().startsWith("""["CLOSED","restricted:"""), refused.lines.toString()) + assertTrue(refused.lines.single().startsWith("""["CLOSED","q","restricted:"""), refused.lines.toString()) val optingIn = MemoryRelay(backend, { @@ -359,14 +376,14 @@ class HttpRelayHandlerTest { override suspend fun authorizeTransport(pubkey: HexKey): String? = null } }) - val signed = HttpRelayHandler(optingIn, origins = { listOf(origin) }).ask(HttpRelayCommand.REQ, body, token(HttpRelayCommand.REQ, body)) + val signed = HttpRelayHandler(optingIn, origins = { listOf(origin) }).ask(frame, token(frame)) 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)}"}""") + val answer = HttpRelayHandler(limited, origins = { listOf(origin) }).ask(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()) } @@ -374,13 +391,12 @@ class HttpRelayHandlerTest { @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()) + val answer = handler().ask(publish(note("中".repeat(1_500)))) assertEquals(200, answer.status, answer.lines.toString()) } @Test - fun aBackendFailureIsA500Line() { + fun aBackendFailureIsA500OkFalse() { val failing = object : SessionBackend by backend { override suspend fun submit( @@ -388,17 +404,18 @@ class HttpRelayHandlerTest { onComplete: (IEventStore.InsertOutcome) -> Unit, ): Unit = error("db is down") } - val answer = HttpRelayHandler(MemoryRelay(failing, { VerifyPolicy }), origins = { listOf(origin) }).ask(HttpRelayCommand.EVENT, note("lost").toJson()) + val lost = note("lost") + val answer = HttpRelayHandler(MemoryRelay(failing, { VerifyPolicy }), origins = { listOf(origin) }).ask(publish(lost)) assertEquals(500, answer.status, answer.lines.toString()) - assertTrue(answer.lines.single().startsWith("""["CLOSED","error:"""), answer.lines.toString()) + assertTrue(answer.lines.single().startsWith("""["OK","${lost.id}",false,"error:"""), answer.lines.toString()) } @Test fun anInfiniteDeadlineStillStreams() { backend.events += note("forever") - val answer = handler(deadline = Duration.INFINITE).ask(HttpRelayCommand.REQ, """{"kinds":[1]}""") + val answer = handler(deadline = Duration.INFINITE).ask(req("""{"kinds":[1]}""")) assertEquals(200, answer.status) - assertEquals("""["EOSE"]""", answer.lines.last()) + assertEquals("""["EOSE","q"]""", answer.lines.last()) } @Test @@ -414,37 +431,27 @@ class HttpRelayHandlerTest { override suspend fun stream(lines: suspend HttpRelayLines.() -> Unit) = error("single") } assertFailsWith { - runBlocking { h.handle(HttpRelayRequest(HttpRelayCommand.COUNT, null, """{"kinds":[1]}""".encodeToByteArray()), stalled) } + runBlocking { h.handle(HttpRelayRequest(null, count("""{"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()) + val answer = handler().ask(req("""{"kinds":[1]}""") + publish(smuggled)) + assertTrue(answer.lines.last().let { it == """["EOSE","q"]""" || 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 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) + val frame = req("""{"kinds":[1]}""") + assertEquals(200, h.ask(frame, token(frame, onion)).status) + assertEquals(200, h.ask(frame, token(frame, "http://relayxyz.onion")).status, "with or without the trailing slash") + assertEquals(200, h.ask(frame, token(frame, "$origin/")).status) + assertEquals(401, h.ask(frame, token(frame, "https://elsewhere.example")).status) } private companion object {