mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
fix: NIP-FE audit — slow bodies, malformed Content-Type, body buffers, one parse, faster gzip
Bugs (reproduced on a live geode before fixing, now regression tests): - A body trickled in forever held its admission slot forever: a client that declared 100 bytes and sent 5 still had its slot after 20 s, so a few addresses could fill the global cap. The body read is now bounded by [http].body_timeout_seconds (10 s): 408, Connection: close. - A malformed Content-Type made the NIP-86 check throw, answering 500. It is now compared as text, and a command needs no type at all. Performance: - Every request allocated a cap-sized body buffer (512 KiB) for what is usually a 50-byte frame. The buffer is now sized to Content-Length, or grows from 4 KiB for a chunked body. - The session parsed each body again after the handler had. RelaySession.receive(text, parsed) runs the raw-text policies on the text and dispatches the parsed command. - Gzip at level 1: 42% of raw against 39% at the default 6 on Nostr-shaped events, for about 2.5x less CPU. Compressed bytes go to the socket from their own buffer instead of a copy. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01RbNrTdV2e7kW5S9tkMoPgh
This commit is contained in:
@@ -182,6 +182,8 @@ enabled = true
|
||||
# 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
|
||||
|
||||
@@ -242,6 +242,8 @@ data class StaticConfig(
|
||||
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. */
|
||||
@@ -270,6 +272,7 @@ data class StaticConfig(
|
||||
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,
|
||||
@@ -357,6 +360,7 @@ data class StaticConfig(
|
||||
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}" }
|
||||
|
||||
@@ -28,6 +28,8 @@ 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,
|
||||
@@ -36,13 +38,17 @@ internal suspend fun readBoundedBody(
|
||||
val declared = call.request.headers[HttpHeaders.ContentLength]?.toLongOrNull()
|
||||
if (declared != null && declared > cap) return null
|
||||
val ch = call.receiveChannel()
|
||||
val buf = ByteArray(cap + 1)
|
||||
// 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 buf.copyOfRange(0, pos)
|
||||
return if (pos == buf.size) buf else buf.copyOf(pos)
|
||||
}
|
||||
|
||||
private const val INITIAL_BODY_BUFFER = 4 * 1024
|
||||
|
||||
@@ -23,6 +23,7 @@ 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
|
||||
|
||||
/**
|
||||
@@ -34,8 +35,16 @@ import java.util.zip.GZIPOutputStream
|
||||
internal class GzipLines(
|
||||
private val out: ByteWriteChannel,
|
||||
) {
|
||||
private val compressed = ByteArrayOutputStream(BUFFER)
|
||||
private val gzip = GZIPOutputStream(compressed, BUFFER, true)
|
||||
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) {
|
||||
@@ -63,10 +72,15 @@ internal class GzipLines(
|
||||
|
||||
private suspend fun drain() {
|
||||
if (compressed.size() == 0) return
|
||||
out.writeFully(compressed.toByteArray())
|
||||
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())
|
||||
|
||||
@@ -23,6 +23,7 @@ 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
|
||||
@@ -35,6 +36,8 @@ data class HttpCommandSettings(
|
||||
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. */
|
||||
|
||||
@@ -31,13 +31,14 @@ 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.contentType
|
||||
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
|
||||
@@ -71,11 +72,11 @@ internal class NipFEHttpRoute(
|
||||
* 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 =
|
||||
!call.request
|
||||
.contentType()
|
||||
.withoutParameters()
|
||||
.match(NIP86)
|
||||
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) {
|
||||
@@ -87,13 +88,22 @@ internal class NipFEHttpRoute(
|
||||
|
||||
val verdict =
|
||||
admission.admit(clientOf(call)) {
|
||||
// Bounded in time too: a body trickled in forever would hold this admission slot forever.
|
||||
val body =
|
||||
readBoundedBody(call, bodyCap)
|
||||
?: return@admit respondLine(
|
||||
try {
|
||||
withTimeout(settings.bodyTimeout) { readBoundedBody(call, bodyCap) }
|
||||
} catch (_: TimeoutCancellationException) {
|
||||
call.response.header(HttpHeaders.Connection, "close")
|
||||
return@admit respondLine(
|
||||
call,
|
||||
HttpRelayStatus.PAYLOAD_TOO_LARGE,
|
||||
HttpRelayHandler.notice(MachineReadablePrefix.INVALID.format("the command exceeds $bodyCap bytes")),
|
||||
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) {
|
||||
@@ -213,7 +223,7 @@ internal class NipFEHttpRoute(
|
||||
|
||||
companion object {
|
||||
val NDJSON = ContentType("application", "x-ndjson")
|
||||
private val NIP86 = ContentType.parse(Nip86HttpHandler.CONTENT_TYPE)
|
||||
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
|
||||
|
||||
@@ -58,6 +58,7 @@ 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
|
||||
@@ -66,6 +67,7 @@ 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
|
||||
|
||||
/**
|
||||
@@ -389,6 +391,65 @@ class NipFEHttpTest {
|
||||
}
|
||||
}
|
||||
|
||||
/** 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()
|
||||
|
||||
+19
-4
@@ -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()) {
|
||||
|
||||
+4
-4
@@ -121,9 +121,9 @@ class HttpRelayHandler(
|
||||
// Characters, as the engine's own limit counts them.
|
||||
if (max != null && text.length > max) return response.single(HttpRelayStatus.PAYLOAD_TOO_LARGE, notice(tooLarge))
|
||||
|
||||
// Read here as well as in the engine, to refuse the commands HTTP does not carry and to
|
||||
// know what ends the answer and how to refuse it. The engine still parses the text itself,
|
||||
// as socket text, so the policies that judge raw frames run.
|
||||
// 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)
|
||||
@@ -188,7 +188,7 @@ class HttpRelayHandler(
|
||||
try {
|
||||
server.serve(sink) { session ->
|
||||
val refused = signedIn?.let { session.authenticateByTransport(it) }
|
||||
if (refused != null) fail(refused) else session.receive(text)
|
||||
if (refused != null) fail(refused) else session.receive(text, cmd)
|
||||
ended.await()
|
||||
}
|
||||
} catch (e: CancellationException) {
|
||||
|
||||
Reference in New Issue
Block a user