Merge pull request #4231 from vitorpamplona/claude/elegant-bardeen-kt25tq

Add NIP-FE (HTTP relay commands) support to geode
This commit is contained in:
Vitor Pamplona
2026-09-27 14:36:56 -04:00
committed by GitHub
25 changed files with 2032 additions and 271 deletions
+30 -4
View File
@@ -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","<id>",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
+39 -1
View File
@@ -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,
@@ -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) {
@@ -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<String>) {
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<IRelayPolicy>()
if (requireAuth) {
pieces += FullAuthPolicy(advertisedUrl)
pieces += SignInPolicy(advertisedUrl)
} else if (optionalAuth) {
pieces += OptionalAuthPolicy(advertisedUrl)
pieces += OptionalSignInPolicy(advertisedUrl)
}
config.options.reject_future_seconds?.let { secs ->
@@ -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<String> =
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 =
@@ -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
}
@@ -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<MirrorSection> = 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<Int>? = null,
/** NIP numbers, and the hex-named NIPs as strings: `[1, 11, 42, "FE"]`. */
val supported_nips: List<String>? = null,
val privacy_policy: String? = null,
val terms_of_service: String? = null,
val relay_countries: List<String>? = 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<String> = 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<String> = 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<String> = emptyList(),
val pubkey_blacklist: List<String> = 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<StaticConfig>(toml)
@@ -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
@@ -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())
}
}
@@ -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<String, Int>()
/** 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 }
}
}
@@ -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<NormalizedRelayUrl> = 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<String> = 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",
)
@@ -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
@@ -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
}
}
@@ -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<Pair<KtorRelay, RelayEngine>>()
@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<OkMessage>(client().publish(relay, it).last).success) }
val got = mutableListOf<Event>()
val answer = client().req(relay, listOf(Filter(kinds = listOf(TextNoteEvent.KIND))), onEvent = got::add)
assertEquals(200, answer.status)
assertTrue(answer.complete)
assertIs<EoseMessage>(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<CountMessage>(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<ReqCmd> = 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<OkMessage>(published.last).success)
val got = mutableListOf<Event>()
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<Unit>): (IEventStore) -> IEventStore =
{ real ->
object : IEventStore by real {
override suspend fun rawQuery(
filters: List<Filter>,
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<Unit>()
val relay = start(store = holdingEose(release))
val held = alice.sign<Event>(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<Unit>()
val relay = start(store = holdingEose(release))
val held = alice.sign<Event>(TimeUtils.now(), HELD_KIND, emptyArray(), "held")
client().publish(relay, held)
val first = CompletableDeferred<Event>()
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
}
}
@@ -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<IllegalArgumentException> { StaticConfig.fromToml("[http]\ndeadline_seconds = 0").validate() }
assertFailsWith<IllegalArgumentException> { StaticConfig.fromToml("[http]\nmax_requests_per_client = -1").validate() }
}
@Test
fun loadsTheBundledExampleConfigCleanly() {
// The example file lives at the module root so operators have a
@@ -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<Unit>()
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<Unit>()
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") {})
}
}
@@ -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()) {
@@ -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)
}
@@ -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<Filter>,
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<Filter>,
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)
}
}
@@ -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)
}
}
}
@@ -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<String>,
/** 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)
}
@@ -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,
)
}
@@ -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<EoseMessage>(reader.last)
}
@Test
fun aReqWithItsTailMissingIsIncomplete() {
val reader = HttpRelayAnswerReader(HttpRelayCommand.REQ, 200)
val message = reader.read("""["EVENT","q",$event]""")
assertIs<EventMessage>(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<CountMessage>(count.last).result.count)
val ok = read(HttpRelayCommand.EVENT, 200, """["OK","abc",true,""]""")
assertTrue(ok.complete)
assertTrue(assertIs<OkMessage>(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<ClosedMessage>(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(""))
}
}
@@ -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()
}
}
@@ -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<Event>(at, 1, emptyArray(), content)
private fun token(
command: HttpRelayCommand,
body: String,
) = alice.sign(HTTPAuthorizationEvent.build(origin + command.path, "POST", body.encodeToByteArray(), System.currentTimeMillis() / 1000) {}).toAuthToken()
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<HttpRelayReaderStalled> {
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<HttpRelayReaderStalled> {
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 {