feat: gzip NIP-FE answers with a sync flush; pluggable HTTP transport

geode:
- Streamed answers are gzipped when the client sends
  Accept-Encoding: gzip (not q=0), with Vary: Accept-Encoding. GzipLines
  sync-flushes at every flush the handler makes, so lines still arrive
  as they are found; compressed bytes also move to the socket past
  16 KiB so a long burst keeps backpressure. JDK zip only, no new
  dependency. One-line answers stay plain. [http].gzip turns it off.
- Tests hold EOSE in the store after the first event and require that
  event to arrive, inflated, before EOSE is released: read raw through
  okio's GzipSource and through HttpRelayClient. Both time out when the
  sync flush is turned off.

quartz:
- HttpRelayClient moves to commonMain over an HttpRelayTransport
  interface, as NostrClient sits over a WebsocketBuilder.
  OkHttpRelayTransport (jvmAndroid) picks the OkHttpClient per relay,
  like BasicOkHttpWebSocket.Builder, so .onion relays can go via Tor.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01RbNrTdV2e7kW5S9tkMoPgh
This commit is contained in:
Claude
2026-09-27 17:07:14 +00:00
parent bb92054f65
commit 8f3e74d229
10 changed files with 408 additions and 78 deletions
+5 -2
View File
@@ -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
+4
View File
@@ -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.
@@ -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(),
@@ -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())
}
}
@@ -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. */
@@ -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 {
@@ -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<OkMessage>(published.last).success)
val got = mutableListOf<Event>()
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<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())
}
}
}
@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
}
}
@@ -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)
}
}
@@ -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,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()
}
}