diff --git a/geode/README.md b/geode/README.md index a32af00068..11732aaad9 100644 --- a/geode/README.md +++ b/geode/README.md @@ -108,7 +108,7 @@ 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 -d '["REQ","q",{"kinds":[1],"limit":2}]' http://localhost:7447/ +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"] @@ -121,7 +121,10 @@ A body that does not end on `EOSE`/`CLOSED` (REQ), `COUNT`/`CLOSED` (COUNT) or 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`. -Quartz's `HttpRelayClient` does all of this for JVM/Android clients. +Streamed answers are gzipped for clients that send `Accept-Encoding: gzip` +(`curl --compressed`), sync-flushed so events still arrive as they are found; +`[http].gzip = false` turns it off. Quartz's `HttpRelayClient` does all of this +for clients, over OkHttp via `OkHttpRelayTransport` or any `HttpRelayTransport`. ## Verbs diff --git a/geode/config.example.toml b/geode/config.example.toml index a2a46fe6c7..891ed025dc 100644 --- a/geode/config.example.toml +++ b/geode/config.example.toml @@ -185,6 +185,10 @@ max_requests_per_client = 16 # 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. diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/config/StaticConfig.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/config/StaticConfig.kt index b3d888565c..6212aa7178 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/config/StaticConfig.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/config/StaticConfig.kt @@ -246,6 +246,8 @@ data class StaticConfig( 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, /** @@ -270,6 +272,7 @@ data class StaticConfig( maxPerClient = max_requests_per_client, deadline = deadline_seconds.seconds, maxBodyBytes = max_body_bytes, + compress = gzip, retryAfterSeconds = retry_after_seconds, alternateUrls = alternate_urls.map { it.normalizeRelayUrl() }, trustedProxies = trusted_proxies.toSet(), diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/server/GzipLines.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/server/GzipLines.kt new file mode 100644 index 0000000000..aa4f4c4e22 --- /dev/null +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/server/GzipLines.kt @@ -0,0 +1,74 @@ +/* + * 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.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 = ByteArrayOutputStream(BUFFER) + private val gzip = GZIPOutputStream(compressed, BUFFER, true) + + /** 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 + out.writeFully(compressed.toByteArray()) + compressed.reset() + } + + companion object { + private const val BUFFER = 16 * 1024 + private val NEWLINE = byteArrayOf('\n'.code.toByte()) + } +} diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/server/HttpCommandSettings.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/server/HttpCommandSettings.kt index 5d1017889e..0508511eb7 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/server/HttpCommandSettings.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/server/HttpCommandSettings.kt @@ -39,6 +39,8 @@ data class HttpCommandSettings( 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. */ diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/server/NipFEHttpRoute.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/server/NipFEHttpRoute.kt index 77b7ac51b3..8980c3031c 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/server/NipFEHttpRoute.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/server/NipFEHttpRoute.kt @@ -132,6 +132,27 @@ internal class NipFEHttpRoute( ?.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, @@ -153,19 +174,41 @@ internal class NipFEHttpRoute( frame: String, ) = respondLine(call, status, frame) - /** Chunked: each flush puts what the handler wrote on the wire, and a full socket suspends the writer. */ - override suspend fun stream(lines: suspend HttpRelayLines.() -> Unit) = + /** + * 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 - object : HttpRelayLines { - override suspend fun line(frame: String) { - out.writeStringUtf8(frame) - out.writeStringUtf8("\n") - } + if (gzip) { + val body = GzipLines(out) + try { + object : HttpRelayLines { + override suspend fun line(frame: String) = body.line(frame) - override suspend fun flush() = out.flush() - }.lines() + 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 { diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/NipFEHttpTest.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/NipFEHttpTest.kt index 8650bfb503..0403e7fde4 100644 --- a/geode/src/test/kotlin/com/vitorpamplona/geode/NipFEHttpTest.kt +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/NipFEHttpTest.kt @@ -31,21 +31,34 @@ 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.util.concurrent.TimeUnit import kotlin.test.AfterTest import kotlin.test.Test import kotlin.test.assertEquals @@ -53,6 +66,7 @@ import kotlin.test.assertFalse import kotlin.test.assertIs import kotlin.test.assertTrue import kotlin.time.Duration +import kotlin.time.Duration.Companion.seconds /** * NIP-FE end to end: quartz's [HttpRelayClient] and raw OkHttp requests against a real [KtorRelay], @@ -77,16 +91,18 @@ class NipFEHttpTest { 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 relay = if (policy == null) RelayEngine(url) else RelayEngine(url, policyBuilder = { policy(url) }) + 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(http, signer) + private fun client(signer: NostrSignerInternal? = null) = HttpRelayClient(OkHttpRelayTransport { http }, signer) private suspend fun note(text: String): Event = alice.sign(TextNoteEvent.build(text)) @@ -239,7 +255,7 @@ class NipFEHttpTest { assertTrue(assertIs(published.last).success) val got = mutableListOf() - val read = HttpRelayClient(http, alice, signFirst = true).req(relay, listOf(Filter(ids = listOf(n.id))), onEvent = got::add) + 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 }) @@ -297,6 +313,82 @@ class NipFEHttpTest { } } + /** Answers a filter naming [HELD_KIND] with what is stored, then holds its EOSE until [release]. */ + private fun holdingEose(release: CompletableDeferred): (IEventStore) -> IEventStore = + { real -> + object : IEventStore by real { + override suspend fun rawQuery( + filters: List, + onEach: (RawEvent) -> Unit, + ) { + real.rawQuery(filters, onEach) + if (filters.any { it.kinds?.contains(HELD_KIND) == true }) release.await() + } + } + } + + @Test + fun aGzippedAnswerStillArrivesLineByLine() = + runBlocking { + val release = CompletableDeferred() + val relay = start(store = holdingEose(release)) + val held = alice.sign(TimeUtils.now(), HELD_KIND, emptyArray(), "held") + client().publish(relay, held) + + // Asking for gzip by hand turns off OkHttp's own inflating, so the body is read as sent. + val patient = http.newBuilder().readTimeout(10, TimeUnit.SECONDS).build() + val request = + Request + .Builder() + .url(relay.toHttp()) + .post("""["REQ","q",{"kinds":[$HELD_KIND]}]""".toRequestBody("text/plain".toMediaType())) + .header("Accept-Encoding", "gzip") + .build() + patient.newCall(request).execute().use { response -> + assertEquals("gzip", response.header("Content-Encoding")) + val lines = GzipSource(response.body.source()).buffer() + // The store has not answered EOSE: this line can only be here if the gzip was flushed. + assertEquals("""["EVENT","q",${held.toJson()}]""", lines.readUtf8Line()) + assertFalse(release.isCompleted) + release.complete(Unit) + assertEquals("""["EOSE","q"]""", lines.readUtf8Line()) + assertEquals(null, lines.readUtf8Line()) + } + } + + @Test + fun theClientReadsAGzippedAnswerAsItStreams() = + runBlocking { + val release = CompletableDeferred() + val relay = start(store = holdingEose(release)) + val held = alice.sign(TimeUtils.now(), HELD_KIND, emptyArray(), "held") + client().publish(relay, held) + + val first = CompletableDeferred() + val answer = async(Dispatchers.IO) { client().req(relay, listOf(Filter(kinds = listOf(HELD_KIND))), onEvent = { first.complete(it) }) } + assertEquals(held.id, withTimeout(10.seconds) { first.await() }.id, "the event arrives before EOSE is written") + release.complete(Unit) + assertTrue(answer.await().complete) + } + + @Test + fun withoutGzipOrWithItOffTheBodyIsPlain() { + for ((settings, accept) in listOf(HttpCommandSettings() to "identity", HttpCommandSettings(compress = false) to "gzip")) { + val relay = start(settings = settings) + val request = + Request + .Builder() + .url(relay.toHttp()) + .post("""["REQ","q",{"kinds":[1]}]""".toRequestBody("text/plain".toMediaType())) + .header("Accept-Encoding", accept) + .build() + http.newCall(request).execute().use { response -> + assertEquals(null, response.header("Content-Encoding"), accept) + assertEquals(listOf("""["EOSE","q"]"""), response.lines()) + } + } + } + @Test fun nip11AdvertisesFE() { val relay = start() @@ -321,4 +413,9 @@ class NipFEHttpTest { assertTrue(answer.complete) assertFalse(answer.last is EoseMessage) } + + private companion object { + /** A regular kind (stored), not a text note, so no other test reads it. */ + const val HELD_KIND = 7_777 + } } diff --git a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayClient.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayClient.kt similarity index 61% rename from quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayClient.kt rename to quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayClient.kt index 5eb4bd6365..af7eb13305 100644 --- a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayClient.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayClient.kt @@ -33,22 +33,12 @@ 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 -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 /** - * NIP-FE over OkHttp: one relay command per request, POSTed to the relay's URL as the frame the + * 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. + * 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 @@ -56,7 +46,7 @@ import okio.IOException * relay known to want it. */ class HttpRelayClient( - private val http: OkHttpClient, + private val transport: HttpRelayTransport, private val signer: NostrSigner? = null, private val signFirst: Boolean = false, ) { @@ -92,11 +82,11 @@ class HttpRelayClient( ): HttpRelayAnswer { val command = requireNotNull(HttpRelayCommand.of(cmd)) { "NIP-FE carries REQ, COUNT and EVENT, not ${cmd.label()}" } val url = relay.toHttp() - val bytes = cmd.toJson().encodeToByteArray() - if (signer == null || signFirst) return post(command, url, bytes, token(url, bytes), onMessage, retrying = false) - val first = post(command, url, bytes, null, onMessage, retrying = true) + 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(command, url, bytes, token(url, bytes), onMessage, retrying = false) + return post(relay, command, url, body, token(url, body), onMessage, retrying = false) } private suspend fun token( @@ -105,6 +95,7 @@ class HttpRelayClient( ): String? = signer?.sign(HTTPAuthorizationEvent.build(url, "POST", body))?.toAuthToken() private suspend fun post( + relay: NormalizedRelayUrl, command: HttpRelayCommand, url: String, body: ByteArray, @@ -112,51 +103,25 @@ class HttpRelayClient( onMessage: (Message) -> Unit, /** A signed try follows a 401, so that refusal is not the answer and is not handed on. */ retrying: Boolean, - ): HttpRelayAnswer = - coroutineScope { - val request = - Request - .Builder() - .url(url) - .post(body.toRequestBody(JSON)) - .header("Accept", NDJSON) - .apply { authorization?.let { header("Authorization", it) } } - .build() - val call = http.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 -> - withContext(Dispatchers.IO) { - val reader = HttpRelayAnswerReader(command, response.code) - val deliver = !(retrying && response.code == HttpRelayStatus.UNAUTHORIZED) - val source = response.body.source() - try { - while (true) { - val line = source.readUtf8Line() ?: break - val message = reader.read(line) - if (message != null && deliver) onMessage(message) - } - } catch (_: IOException) { - // The connection dropped mid-answer: what came is what the reader says it is. - } - reader.answer(response.header("Retry-After")) - } - } - } finally { - watcher.cancel() - } - } - - companion object { - const val NDJSON = "application/x-ndjson" - private val JSON = "application/json".toMediaType() + ): HttpRelayAnswer { + var reader: HttpRelayAnswerReader? = null + var retryAfter: String? = null + var deliver = true + transport.post( + relay = relay, + url = url, + body = body, + authorization = authorization, + onStatus = { status, after -> + reader = HttpRelayAnswerReader(command, status) + retryAfter = after + deliver = !(retrying && status == HttpRelayStatus.UNAUTHORIZED) + }, + onLine = { line -> + val message = checkNotNull(reader) { "a line before the status" }.read(line) + if (message != null && deliver) onMessage(message) + }, + ) + return checkNotNull(reader) { "the transport returned without a status" }.answer(retryAfter) } } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayTransport.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayTransport.kt new file mode 100644 index 0000000000..918bb8e921 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/HttpRelayTransport.kt @@ -0,0 +1,47 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.quartz.nipFERelayOverHttp + +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl + +/** + * How [HttpRelayClient] reaches a relay over HTTP: one POST, its response read line by line. The + * NIP-FE side of the exchange (which URL, what body, signing, reading the answer) stays in the + * client; an implementation only moves bytes, the way a + * [com.vitorpamplona.quartz.nip01Core.relay.sockets.WebsocketBuilder] does for [com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient]. + */ +interface HttpRelayTransport { + /** + * POSTs [body] to [url], [relay]'s HTTP URL, with an `Authorization` header when [authorization] + * is not null. Calls [onStatus] once with the status and any `Retry-After`, then [onLine] for each + * body line as it arrives, and returns when the body ends. A connection that drops mid-body + * returns normally: the answer reader tells a cut-off answer from a whole one. Cancelling the + * caller cancels the request. + */ + suspend fun post( + relay: NormalizedRelayUrl, + url: String, + body: ByteArray, + authorization: String?, + onStatus: (status: Int, retryAfter: String?) -> Unit, + onLine: (String) -> Unit, + ) +} diff --git a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/OkHttpRelayTransport.kt b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/OkHttpRelayTransport.kt new file mode 100644 index 0000000000..ee65bda983 --- /dev/null +++ b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nipFERelayOverHttp/OkHttpRelayTransport.kt @@ -0,0 +1,92 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.quartz.nipFERelayOverHttp + +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.awaitCancellation +import kotlinx.coroutines.coroutineScope +import kotlinx.coroutines.launch +import kotlinx.coroutines.withContext +import okhttp3.MediaType.Companion.toMediaType +import okhttp3.OkHttpClient +import okhttp3.Request +import okhttp3.RequestBody.Companion.toRequestBody +import okhttp3.coroutines.executeAsync +import okio.IOException + +/** + * [HttpRelayTransport] over OkHttp. [httpClient] picks the client per relay, as + * [com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.BasicOkHttpWebSocket.Builder] does for + * websockets, so a .onion relay can go through Tor and a clearnet one directly. OkHttp asks for + * gzip and inflates it as the body streams, so a relay's sync-flushed gzip arrives line by line. + */ +class OkHttpRelayTransport( + private val httpClient: (NormalizedRelayUrl) -> OkHttpClient, +) : HttpRelayTransport { + override suspend fun post( + relay: NormalizedRelayUrl, + url: String, + body: ByteArray, + authorization: String?, + onStatus: (status: Int, retryAfter: String?) -> Unit, + onLine: (String) -> Unit, + ) = coroutineScope { + val request = + Request + .Builder() + .url(url) + .post(body.toRequestBody(JSON)) + .header("Accept", NDJSON) + .apply { authorization?.let { header("Authorization", it) } } + .build() + val call = httpClient(relay).newCall(request) + // A blocking read does not see coroutine cancellation; cancelling the call unblocks it. + val watcher = + launch { + try { + awaitCancellation() + } finally { + call.cancel() + } + } + try { + call.executeAsync().use { response -> + onStatus(response.code, response.header("Retry-After")) + withContext(Dispatchers.IO) { + val source = response.body.source() + try { + while (true) onLine(source.readUtf8Line() ?: break) + } catch (_: IOException) { + // The connection dropped mid-answer: what came is what the reader says it is. + } + } + } + } finally { + watcher.cancel() + } + } + + companion object { + const val NDJSON = "application/x-ndjson" + private val JSON = "application/json".toMediaType() + } +}