mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
feat: add production benchmark for the NostrClient receive path
Gated JVM test (-PprodRelayBench=1) that connects to live relays with realistic filters and measures per-relay queue delay, processing time, consumer busy fraction, EOSE latency and duplicate rates, comparing the current inline verification against a parallel verify stage, plus offline single-thread parse/verify ceilings on captured frames. Findings from a first run are written up in quartz/plans/2026-07-02-nostrclient-receiver-perf.md. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_018saXqYfAa3RvSJoDXK591R
This commit is contained in:
@@ -85,6 +85,12 @@ kotlin {
|
||||
tasks.withType<Test>().configureEach {
|
||||
maxHeapSize = "4g"
|
||||
environment("TEST_RESOURCES_ROOT", rootDir)
|
||||
// Opt-in gate for ProductionReceiverBenchmark (opens sockets to live
|
||||
// relays). `-PprodRelayBench=1` reaches the daemon even when an env
|
||||
// var would not survive daemon reuse.
|
||||
(project.findProperty("prodRelayBench") as? String)?.let {
|
||||
environment("PROD_RELAY_BENCH", it)
|
||||
}
|
||||
}
|
||||
|
||||
tasks.withType<KotlinNativeTest>().configureEach {
|
||||
|
||||
@@ -0,0 +1,144 @@
|
||||
# NostrClient receiver performance — review + production measurements
|
||||
|
||||
**Date:** 2026-07-02
|
||||
**Status:** findings + measurement harness; no production code changed yet
|
||||
**Harness:** `quartz/src/jvmTest/.../relay/prodbench/ProductionReceiverBenchmark.kt`
|
||||
(run with `./gradlew :quartz:jvmTest --tests "*.ProductionReceiverBenchmark" -PprodRelayBench=1`)
|
||||
|
||||
## Question
|
||||
|
||||
Is the NostrClient's single-threaded receiver making us significantly slower?
|
||||
And how do our filters actually behave against production relays?
|
||||
|
||||
## What the receiver actually is (review)
|
||||
|
||||
The receive path is single-threaded **per relay connection**, not globally:
|
||||
|
||||
```
|
||||
OkHttp reader thread (per socket)
|
||||
└─ trySendBlocking → Channel(UNLIMITED) [BasicOkHttpWebSocket / app OkHttpWebSocket]
|
||||
└─ ONE consumer coroutine per connection (Dispatchers.IO)
|
||||
├─ OptimizedJsonMapper.fromJsonToMessage [JSON parse] ~32µs/msg (JVM)
|
||||
├─ PoolRequests.onIncomingMessage [spin-lock + state]
|
||||
│ └─ SubscriptionListener.onEvent [assemblers, inline]
|
||||
└─ NostrClient.listeners.forEach [EventCollector, loggers…]
|
||||
└─ LocalCache.justConsume(…, wasVerified = false)
|
||||
└─ dedup by id → justVerify() [Schnorr verify] ~75µs/event (JVM)
|
||||
└─ note creation, index updates, flow emissions
|
||||
```
|
||||
|
||||
Everything downstream of the channel — parse, subscription state machine, every
|
||||
listener, signature verification, and LocalCache consumption — runs serially on
|
||||
that one consumer coroutine. Different relays run on different coroutines
|
||||
(cross-relay parallelism exists), but they contend on shared state: the
|
||||
`PoolRequests` spin lock, `LocalCache`'s maps, and (on Android) the same
|
||||
`Dispatchers.IO` pool as everything else.
|
||||
|
||||
Notably, the **relay-server side already solved this** on the write path:
|
||||
`IngestQueue` batches incoming EVENTs and runs Schnorr verification for a batch
|
||||
in parallel on `Dispatchers.Default` (`parallelVerify`), precisely because "an
|
||||
event's `verify()` is CPU-bound and parallelisable". The client receiver has no
|
||||
equivalent.
|
||||
|
||||
## Production measurements (2026-07-02)
|
||||
|
||||
Environment: JVM 21, 4-core x86 cloud container (server CPU — phones are
|
||||
5–10× slower per core), relays `relay.damus.io`, `nos.lol`, `relay.primal.net`,
|
||||
`nostr.wine`. All scenarios ran after a discarded warmup pass (JIT warm).
|
||||
`INLINE` = verify on the receiver coroutine (what the app does today via
|
||||
`CacheClientConnector`); `PARALLEL` = verify offloaded to a
|
||||
`Dispatchers.Default` pool.
|
||||
|
||||
### Filter behavior
|
||||
|
||||
| Scenario | received | unique | duplicates | notes |
|
||||
|---|---|---|---|---|
|
||||
| kind-1 firehose, limit 500 ×4 relays | 2013 | 1738 | 14% | all relays honored limit; EOSE 0.3–1.0s after REQ |
|
||||
| notifications (p=busy pubkey, kinds 1,6,7,9735) | 1552 | 986 | 36% | primal returned only 52 (thin index for this query) |
|
||||
| metadata burst (kinds 0,10002, 300 authors) | 1146 | 491 | 57% | each profile fetched ~2.3× across 4 relays |
|
||||
|
||||
- Time from `subscribe()` to REQ-on-the-wire is dominated by connect/TLS
|
||||
(~0.7–1.5s cold); local dispatch is negligible.
|
||||
- Event sizes vary wildly by relay: damus kind-1 averaged **8.7 KB**/event
|
||||
(4.4 MB for 500 events) vs ~0.8 KB on nos.lol — parse cost tracks size
|
||||
(p90 proc 1.5ms on damus vs 0.18ms on nos.lol).
|
||||
- Duplicate factor grows with fan-out: the metadata pattern wastes >half the
|
||||
received bytes on duplicates. Dedup happens *before* verify (LocalCache
|
||||
checks the id first), so duplicates cost parse + dedup lookup, not a verify.
|
||||
|
||||
### Receiver queue behavior (the actual question)
|
||||
|
||||
Queue delay = how long a frame sat in the per-connection channel before the
|
||||
consumer picked it up. This is the direct symptom of a saturated receiver.
|
||||
|
||||
kind-1 firehose (large events), per relay:
|
||||
|
||||
| | INLINE p50 / p90 / max | PARALLEL p50 / p90 / max |
|
||||
|---|---|---|
|
||||
| relay.damus.io | 6.8ms / 26.6ms / 77.9ms | 5.4ms / 38.1ms / 41.5ms |
|
||||
| nos.lol | 4.3ms / 11.2ms / 15.9ms | 1.3ms / 3.8ms / 4.5ms |
|
||||
| relay.primal.net | 5.4ms / 24.5ms / 28.1ms | 1.2ms / 1.8ms / 7.7ms |
|
||||
| nostr.wine | 9.4ms / 30.3ms / 34.8ms | 2.4ms / 8.4ms / 20.7ms |
|
||||
|
||||
metadata burst (small events, fastest arrival — damus delivered at 3300–8400 msg/s):
|
||||
|
||||
| | INLINE p50 / max | PARALLEL p50 / max |
|
||||
|---|---|---|
|
||||
| relay.damus.io | 3.0ms / 10.7ms | 0.14ms / 0.48ms |
|
||||
|
||||
Consumer busy fraction peaked at **53%** (nostr.wine, INLINE, 3000 msg/s) on
|
||||
this fast CPU — i.e. a single relay bursting at production speed already
|
||||
consumes half of one core with inline verification, on hardware much faster
|
||||
than a phone.
|
||||
|
||||
### Single-thread ceilings (offline, captured production frames)
|
||||
|
||||
- JSON parse: **31,000 msg/s** (32µs/msg, mixed sizes incl. damus 8.7KB events)
|
||||
- Schnorr verify: **13,200 events/s** sequential (75µs/event);
|
||||
**34,400 events/s** across 4 cores (2.6–3.0× speedup)
|
||||
- Combined parse+verify ceiling for one receiver coroutine: **~9,000 events/s**
|
||||
on this CPU. On a phone core (5–10× slower): **~1,000–2,000 events/s**, and
|
||||
that is *before* LocalCache's note creation, index updates and flow emissions,
|
||||
which the harness does not model and which run on the same coroutine.
|
||||
|
||||
## Answer
|
||||
|
||||
1. **The receiver is not globally single-threaded** — it is one consumer per
|
||||
relay. Cross-relay parallelism already exists. The per-relay serialization
|
||||
is what caps throughput.
|
||||
|
||||
2. **For steady-state browsing the receiver is not the bottleneck.** Queue
|
||||
delays with today's inline design are single-digit-to-tens of ms; connect
|
||||
latency and relay EOSE times (300ms–1.5s) dominate what the user perceives.
|
||||
|
||||
3. **For burst ingest (app start, feed switch, big profile sync) it is a real
|
||||
cost, and it scales linearly with burst size.** A 5,000-event burst from one
|
||||
fast relay serializes ~375ms of parse+verify on this benchmark machine —
|
||||
plausibly **2–4s per relay on a mid-range phone**, during which that relay's
|
||||
EOSE, OKs and live events all sit behind the backlog. Verification is
|
||||
60–70% of that inline cost and is embarrassingly parallel (measured 2.6–3×
|
||||
on 4 cores; phones have 8).
|
||||
|
||||
4. **Offloading verification cuts queue latency 2–5× and raises the per-relay
|
||||
ceiling ~3×** with no protocol or API change — the exact pattern
|
||||
`IngestQueue.parallelVerify` already uses on the server side.
|
||||
|
||||
## Recommendations (in order of value/risk)
|
||||
|
||||
1. **Move Schnorr verification off the receiver coroutine** in the app's
|
||||
consume path (`CacheClientConnector` → `LocalCache.justConsume`): a bounded
|
||||
verify stage on `Dispatchers.Default` (batch + `async` per batch, or a small
|
||||
worker pool), mirroring `IngestQueue`. Ordering constraint: dedup must stay
|
||||
before verify (it already is), and consumers of `justConsume`'s return value
|
||||
need an async-tolerant path.
|
||||
2. **Keep parse on the receiver coroutine** (32µs/msg is cheap and keeps
|
||||
message ordering per subscription), but consider skipping re-serialization
|
||||
work for duplicates (57% of metadata-burst traffic never needs more than an
|
||||
id lookup).
|
||||
3. **Replace the `PoolRequests` spin lock's busy-wait** with a short-critical-
|
||||
section `Mutex`/`synchronized` if profiling on-device shows contention —
|
||||
with dozens of relay consumers on a phone, spinning burns cores the verify
|
||||
pool needs. (Not measurable as a problem in this 4-relay harness.)
|
||||
4. Re-run this harness on-device (the same class compiles for Android
|
||||
instrumentation with minor changes) before/after any fix; the offline
|
||||
ceilings section gives the numbers to compare.
|
||||
+620
@@ -0,0 +1,620 @@
|
||||
/*
|
||||
* 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.nip01Core.relay.prodbench
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper
|
||||
import com.vitorpamplona.quartz.nip01Core.crypto.verify
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.EventCollector
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener
|
||||
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.sockets.WebSocket
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocketListener
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebsocketBuilder
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.async
|
||||
import kotlinx.coroutines.awaitAll
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.channels.Channel
|
||||
import kotlinx.coroutines.channels.trySendBlocking
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.withTimeoutOrNull
|
||||
import okhttp3.OkHttpClient
|
||||
import okhttp3.Request
|
||||
import okhttp3.Response
|
||||
import java.util.Collections
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
import java.util.concurrent.TimeUnit
|
||||
import java.util.concurrent.atomic.AtomicLong
|
||||
import kotlin.test.Test
|
||||
|
||||
/**
|
||||
* Production measurement harness for the NostrClient receive path.
|
||||
*
|
||||
* The client's receiver is "single threaded" PER RELAY: OkHttp's socket-reader
|
||||
* thread enqueues each raw frame into an unbounded Channel, and one consumer
|
||||
* coroutine per connection drains it, doing — inline, serially — the JSON
|
||||
* parse, the PoolRequests state machine, every SubscriptionListener callback,
|
||||
* and every NostrClient connection listener (in the app that includes
|
||||
* LocalCache.justConsume, i.e. Schnorr signature verification of every new
|
||||
* event). This harness quantifies what that costs against real relays:
|
||||
*
|
||||
* - per-relay throughput (msgs/s, bytes) and time-to-EOSE for realistic filters
|
||||
* - queue delay: how long frames sit in the channel before the consumer
|
||||
* gets to them (the direct symptom of a saturated single consumer)
|
||||
* - processing time per message (parse + dispatch + verify when inline)
|
||||
* - INLINE vs PARALLEL verification (the same strategy the relay-server's
|
||||
* IngestQueue already uses on the write path)
|
||||
* - offline single-thread ceilings for parse and verify on the captured
|
||||
* production frames
|
||||
*
|
||||
* This test opens live sockets to public relays, so it is gated: it no-ops
|
||||
* unless the environment variable PROD_RELAY_BENCH is set.
|
||||
*
|
||||
* Run with:
|
||||
* PROD_RELAY_BENCH=1 ./gradlew :quartz:jvmTest --tests "*.ProductionReceiverBenchmark"
|
||||
*/
|
||||
class ProductionReceiverBenchmark {
|
||||
companion object {
|
||||
val RELAYS =
|
||||
listOf(
|
||||
"wss://relay.damus.io",
|
||||
"wss://nos.lol",
|
||||
"wss://relay.primal.net",
|
||||
"wss://nostr.wine",
|
||||
)
|
||||
|
||||
// fiatjaf — a pubkey with heavy notification traffic on public relays.
|
||||
const val BUSY_PUBKEY = "3bf0c63fcb93463407af97a5e5ee64fa883d107ef9e558472c4eb9aaaefa459d"
|
||||
|
||||
const val EOSE_TIMEOUT_MS = 30_000L
|
||||
const val LIVE_WINDOW_MS = 2_000L
|
||||
const val DRAIN_TIMEOUT_MS = 30_000L
|
||||
}
|
||||
|
||||
enum class VerifyMode { INLINE, PARALLEL }
|
||||
|
||||
// -----------------------------------------------------------------
|
||||
// Instrumented socket: BasicOkHttpWebSocket + timestamps around the
|
||||
// per-connection channel so queue delay and processing time are visible.
|
||||
// -----------------------------------------------------------------
|
||||
|
||||
class RelayIngestMetrics(
|
||||
val url: NormalizedRelayUrl,
|
||||
) {
|
||||
val queueDelayNanos = Collections.synchronizedList(ArrayList<Long>(8192))
|
||||
val procNanos = Collections.synchronizedList(ArrayList<Long>(8192))
|
||||
val msgCount = AtomicLong(0)
|
||||
val byteCount = AtomicLong(0)
|
||||
|
||||
@Volatile var firstMsgAtNanos = 0L
|
||||
|
||||
@Volatile var lastMsgAtNanos = 0L
|
||||
|
||||
fun record(
|
||||
queueDelay: Long,
|
||||
proc: Long,
|
||||
bytes: Int,
|
||||
now: Long,
|
||||
) {
|
||||
queueDelayNanos.add(queueDelay)
|
||||
procNanos.add(proc)
|
||||
msgCount.incrementAndGet()
|
||||
byteCount.addAndGet(bytes.toLong())
|
||||
if (firstMsgAtNanos == 0L) firstMsgAtNanos = now
|
||||
lastMsgAtNanos = now
|
||||
}
|
||||
}
|
||||
|
||||
class InstrumentedOkHttpWebSocket(
|
||||
val url: NormalizedRelayUrl,
|
||||
val httpClient: OkHttpClient,
|
||||
val out: WebSocketListener,
|
||||
val metrics: RelayIngestMetrics,
|
||||
val frameSink: ((String) -> Unit)?,
|
||||
) : WebSocket {
|
||||
class Frame(
|
||||
val text: String,
|
||||
val enqueuedAtNanos: Long,
|
||||
)
|
||||
|
||||
private var socket: okhttp3.WebSocket? = null
|
||||
|
||||
override fun needsReconnect() = socket == null
|
||||
|
||||
override fun connect() {
|
||||
val request = Request.Builder().url(url.url).build()
|
||||
|
||||
val listener =
|
||||
object : okhttp3.WebSocketListener() {
|
||||
val scope = CoroutineScope(Dispatchers.IO + SupervisorJob())
|
||||
val incoming: Channel<Frame> = Channel(Channel.UNLIMITED)
|
||||
val job =
|
||||
scope.launch {
|
||||
// Mirrors BasicOkHttpWebSocket/OkHttpWebSocket: ONE
|
||||
// consumer per connection; everything downstream of
|
||||
// out.onMessage runs serially right here.
|
||||
for (frame in incoming) {
|
||||
val start = System.nanoTime()
|
||||
frameSink?.invoke(frame.text)
|
||||
out.onMessage(frame.text)
|
||||
val end = System.nanoTime()
|
||||
metrics.record(
|
||||
queueDelay = start - frame.enqueuedAtNanos,
|
||||
proc = end - start,
|
||||
bytes = frame.text.length,
|
||||
now = end,
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
override fun onOpen(
|
||||
webSocket: okhttp3.WebSocket,
|
||||
response: Response,
|
||||
) = out.onOpen(
|
||||
(response.receivedResponseAtMillis - response.sentRequestAtMillis).toInt(),
|
||||
response.headers["Sec-WebSocket-Extensions"]?.contains("permessage-deflate") ?: false,
|
||||
)
|
||||
|
||||
override fun onMessage(
|
||||
webSocket: okhttp3.WebSocket,
|
||||
text: String,
|
||||
) {
|
||||
incoming.trySendBlocking(Frame(text, System.nanoTime()))
|
||||
}
|
||||
|
||||
override fun onClosed(
|
||||
webSocket: okhttp3.WebSocket,
|
||||
code: Int,
|
||||
reason: String,
|
||||
) {
|
||||
incoming.close()
|
||||
job.cancel()
|
||||
scope.cancel()
|
||||
out.onClosed(code, reason)
|
||||
}
|
||||
|
||||
override fun onFailure(
|
||||
webSocket: okhttp3.WebSocket,
|
||||
t: Throwable,
|
||||
response: Response?,
|
||||
) {
|
||||
incoming.close()
|
||||
job.cancel()
|
||||
scope.cancel()
|
||||
out.onFailure(t, response?.code, response?.message)
|
||||
}
|
||||
}
|
||||
|
||||
socket = httpClient.newWebSocket(request, listener)
|
||||
}
|
||||
|
||||
override fun disconnect() {
|
||||
socket?.cancel()
|
||||
socket = null
|
||||
}
|
||||
|
||||
override fun send(msg: String): Boolean = socket?.send(msg) ?: false
|
||||
}
|
||||
|
||||
class InstrumentedBuilder(
|
||||
val httpClient: OkHttpClient,
|
||||
val metricsFor: (NormalizedRelayUrl) -> RelayIngestMetrics,
|
||||
val frameSink: ((String) -> Unit)?,
|
||||
) : WebsocketBuilder {
|
||||
override fun build(
|
||||
url: NormalizedRelayUrl,
|
||||
out: WebSocketListener,
|
||||
) = InstrumentedOkHttpWebSocket(url, httpClient, out, metricsFor(url), frameSink)
|
||||
}
|
||||
|
||||
// -----------------------------------------------------------------
|
||||
// Verification sink: mimics what CacheClientConnector/LocalCache do
|
||||
// with every event (dedup by id, then Schnorr-verify new ones),
|
||||
// either inline on the receiver coroutine (current app behavior) or
|
||||
// offloaded to a parallel Default-dispatcher pool.
|
||||
// -----------------------------------------------------------------
|
||||
|
||||
class VerifyingSink(
|
||||
val mode: VerifyMode,
|
||||
scope: CoroutineScope,
|
||||
) {
|
||||
private val seen: MutableSet<String> = ConcurrentHashMap.newKeySet()
|
||||
val received = AtomicLong(0)
|
||||
val unique = AtomicLong(0)
|
||||
val verified = AtomicLong(0)
|
||||
val invalid = AtomicLong(0)
|
||||
val verifyNanosTotal = AtomicLong(0)
|
||||
val perRelayEvents = ConcurrentHashMap<NormalizedRelayUrl, AtomicLong>()
|
||||
|
||||
private val queue: Channel<Event> = Channel(Channel.UNLIMITED)
|
||||
private val workers =
|
||||
if (mode == VerifyMode.PARALLEL) {
|
||||
List(Runtime.getRuntime().availableProcessors()) {
|
||||
scope.launch(Dispatchers.Default) {
|
||||
for (event in queue) doVerify(event)
|
||||
}
|
||||
}
|
||||
} else {
|
||||
emptyList()
|
||||
}
|
||||
|
||||
fun consume(
|
||||
event: Event,
|
||||
relay: NormalizedRelayUrl,
|
||||
) {
|
||||
received.incrementAndGet()
|
||||
perRelayEvents.getOrPut(relay) { AtomicLong(0) }.incrementAndGet()
|
||||
if (seen.add(event.id)) {
|
||||
unique.incrementAndGet()
|
||||
when (mode) {
|
||||
VerifyMode.INLINE -> doVerify(event)
|
||||
VerifyMode.PARALLEL -> queue.trySend(event)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun doVerify(event: Event) {
|
||||
val start = System.nanoTime()
|
||||
val ok = event.verify()
|
||||
verifyNanosTotal.addAndGet(System.nanoTime() - start)
|
||||
if (ok) verified.incrementAndGet() else invalid.incrementAndGet()
|
||||
}
|
||||
|
||||
suspend fun awaitDrained(timeoutMs: Long): Boolean =
|
||||
withTimeoutOrNull(timeoutMs) {
|
||||
while (verified.get() + invalid.get() < unique.get()) delay(20)
|
||||
true
|
||||
} ?: false
|
||||
|
||||
fun close() {
|
||||
queue.close()
|
||||
}
|
||||
}
|
||||
|
||||
// -----------------------------------------------------------------
|
||||
// Scenario runner
|
||||
// -----------------------------------------------------------------
|
||||
|
||||
class ScenarioReport(
|
||||
val capturedEvents: List<Event>,
|
||||
)
|
||||
|
||||
private fun percentileLine(nanos: List<Long>): String {
|
||||
if (nanos.isEmpty()) return "n=0"
|
||||
val sorted = nanos.sorted()
|
||||
|
||||
fun p(q: Double) = sorted[((sorted.size - 1) * q).toInt()]
|
||||
|
||||
fun fmt(n: Long) =
|
||||
when {
|
||||
n >= 1_000_000 -> "%.1fms".format(n / 1_000_000.0)
|
||||
n >= 1_000 -> "%.1fµs".format(n / 1_000.0)
|
||||
else -> "${n}ns"
|
||||
}
|
||||
return "n=${sorted.size} p50=${fmt(p(0.5))} p90=${fmt(p(0.9))} p99=${fmt(p(0.99))} max=${fmt(sorted.last())}"
|
||||
}
|
||||
|
||||
private suspend fun runScenario(
|
||||
name: String,
|
||||
httpClient: OkHttpClient,
|
||||
relays: List<NormalizedRelayUrl>,
|
||||
filters: Map<NormalizedRelayUrl, List<Filter>>,
|
||||
mode: VerifyMode,
|
||||
captureFrames: MutableList<String>? = null,
|
||||
quiet: Boolean = false,
|
||||
): ScenarioReport {
|
||||
if (!quiet) println("\n=== SCENARIO: $name [verify=$mode] ===")
|
||||
|
||||
val metrics = ConcurrentHashMap<NormalizedRelayUrl, RelayIngestMetrics>()
|
||||
val frameSink: ((String) -> Unit)? =
|
||||
captureFrames?.let { list -> { text: String -> if (list.size < 30_000) list.add(text) } }
|
||||
|
||||
val builder =
|
||||
InstrumentedBuilder(
|
||||
httpClient,
|
||||
{ url -> metrics.getOrPut(url) { RelayIngestMetrics(url) } },
|
||||
frameSink,
|
||||
)
|
||||
|
||||
val scenarioScope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
|
||||
val sink = VerifyingSink(mode, scenarioScope)
|
||||
val capturedEvents = Collections.synchronizedList(ArrayList<Event>(8192))
|
||||
|
||||
val client = NostrClient(builder)
|
||||
val collector =
|
||||
EventCollector(client) { event, relay ->
|
||||
sink.consume(event, relay.url)
|
||||
if (capturedEvents.size < 20_000) capturedEvents.add(event)
|
||||
}
|
||||
|
||||
val startNanos = System.nanoTime()
|
||||
val reqSentAt = ConcurrentHashMap<String, Long>()
|
||||
val firstEventAt = ConcurrentHashMap<NormalizedRelayUrl, Long>()
|
||||
val eoseAt = ConcurrentHashMap<NormalizedRelayUrl, Long>()
|
||||
val cannotConnect = ConcurrentHashMap<NormalizedRelayUrl, String>()
|
||||
val eventsBeforeEose = ConcurrentHashMap<NormalizedRelayUrl, AtomicLong>()
|
||||
val eventsAfterEose = ConcurrentHashMap<NormalizedRelayUrl, AtomicLong>()
|
||||
|
||||
val listener =
|
||||
object : SubscriptionListener {
|
||||
override fun onSubscriptionStarted(
|
||||
relay: String,
|
||||
forFilters: List<Filter>,
|
||||
) {
|
||||
reqSentAt.putIfAbsent(relay, System.nanoTime() - startNanos)
|
||||
}
|
||||
|
||||
override fun onEvent(
|
||||
event: Event,
|
||||
isLive: Boolean,
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
firstEventAt.putIfAbsent(relay, System.nanoTime() - startNanos)
|
||||
val counter = if (eoseAt.containsKey(relay)) eventsAfterEose else eventsBeforeEose
|
||||
counter.getOrPut(relay) { AtomicLong(0) }.incrementAndGet()
|
||||
}
|
||||
|
||||
override fun onEose(
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
eoseAt.putIfAbsent(relay, System.nanoTime() - startNanos)
|
||||
}
|
||||
|
||||
override fun onCannotConnect(
|
||||
relay: NormalizedRelayUrl,
|
||||
message: String,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
cannotConnect.putIfAbsent(relay, message)
|
||||
}
|
||||
}
|
||||
|
||||
val subId = "bench-" + System.nanoTime().toString(16)
|
||||
client.subscribe(subId, filters, listener)
|
||||
|
||||
val allEosed =
|
||||
withTimeoutOrNull(EOSE_TIMEOUT_MS) {
|
||||
while (!eoseAt.keys.containsAll(relays) && cannotConnect.keys.size + eoseAt.keys.size < relays.size) delay(50)
|
||||
true
|
||||
} ?: false
|
||||
|
||||
// watch a short live window after EOSE, like the app does
|
||||
delay(LIVE_WINDOW_MS)
|
||||
|
||||
client.unsubscribe(subId)
|
||||
|
||||
val drainStart = System.nanoTime()
|
||||
val drained = sink.awaitDrained(DRAIN_TIMEOUT_MS)
|
||||
val drainNanos = System.nanoTime() - drainStart
|
||||
|
||||
sink.close()
|
||||
collector.destroy()
|
||||
client.close()
|
||||
scenarioScope.cancel()
|
||||
|
||||
// ---- report ----
|
||||
if (quiet) return ScenarioReport(ArrayList(capturedEvents).distinctBy { it.id })
|
||||
|
||||
if (!allEosed) println("!! not all relays reached EOSE within ${EOSE_TIMEOUT_MS}ms")
|
||||
cannotConnect.forEach { (relay, msg) -> println("!! cannot connect ${relay.url}: $msg") }
|
||||
|
||||
for (relay in relays) {
|
||||
val m = metrics[relay] ?: continue
|
||||
val req = reqSentAt[relay.url]?.let { "%.0fms".format(it / 1e6) } ?: "-"
|
||||
val ttfb = firstEventAt[relay]?.let { "%.0fms".format(it / 1e6) } ?: "-"
|
||||
val eose = eoseAt[relay]?.let { "%.0fms".format(it / 1e6) } ?: "-"
|
||||
val activeNanos = (m.lastMsgAtNanos - m.firstMsgAtNanos).coerceAtLeast(1)
|
||||
val rate = m.msgCount.get() * 1e9 / activeNanos
|
||||
val busyPct = m.procNanos.toList().sum() * 100.0 / activeNanos
|
||||
println(" ${relay.url}")
|
||||
println(
|
||||
" req=$req firstEvent=$ttfb eose=$eose msgs=${m.msgCount.get()} bytes=${m.byteCount.get()} rate=%.0f msg/s consumerBusy=%.0f%%"
|
||||
.format(rate, busyPct),
|
||||
)
|
||||
println(" events: pre-EOSE=${eventsBeforeEose[relay]?.get() ?: 0} post-EOSE=${eventsAfterEose[relay]?.get() ?: 0}")
|
||||
println(" queueDelay: ${percentileLine(m.queueDelayNanos.toList())}")
|
||||
println(" procTime: ${percentileLine(m.procNanos.toList())}")
|
||||
}
|
||||
val uniq = sink.unique.get()
|
||||
val recv = sink.received.get()
|
||||
println(" TOTAL: received=$recv unique=$uniq duplicates=${recv - uniq} verified=${sink.verified.get()} invalid=${sink.invalid.get()}")
|
||||
if (uniq > 0) {
|
||||
println(
|
||||
" verify: total=%.1fms avg=%.1fµs/event drainAfterUnsub=%.1fms drained=%s"
|
||||
.format(
|
||||
sink.verifyNanosTotal.get() / 1e6,
|
||||
sink.verifyNanosTotal.get() / 1e3 / uniq,
|
||||
drainNanos / 1e6,
|
||||
drained,
|
||||
),
|
||||
)
|
||||
}
|
||||
|
||||
return ScenarioReport(ArrayList(capturedEvents).distinctBy { it.id })
|
||||
}
|
||||
|
||||
// -----------------------------------------------------------------
|
||||
// Offline microbenchmarks on captured production frames: what a single
|
||||
// receiver coroutine can do per second, with no network in the way.
|
||||
// -----------------------------------------------------------------
|
||||
|
||||
private fun offlineParseBench(frames: List<String>) {
|
||||
if (frames.isEmpty()) return
|
||||
println("\n=== OFFLINE: single-thread JSON parse of ${frames.size} captured frames ===")
|
||||
// warmup
|
||||
repeat(2) { frames.forEach { runCatching { OptimizedJsonMapper.fromJsonToMessage(it) } } }
|
||||
val start = System.nanoTime()
|
||||
var parsed = 0
|
||||
frames.forEach { if (runCatching { OptimizedJsonMapper.fromJsonToMessage(it) }.isSuccess) parsed++ }
|
||||
val nanos = System.nanoTime() - start
|
||||
println(
|
||||
" parsed=$parsed in %.1fms -> %.0f msg/s (%.1fµs/msg)"
|
||||
.format(nanos / 1e6, parsed * 1e9 / nanos, nanos / 1e3 / parsed),
|
||||
)
|
||||
}
|
||||
|
||||
private fun offlineVerifyBench(events: List<Event>) {
|
||||
if (events.isEmpty()) return
|
||||
val sample = events.take(3000)
|
||||
println("\n=== OFFLINE: Schnorr verify of ${sample.size} unique production events ===")
|
||||
|
||||
// warmup (JIT + secp tables)
|
||||
sample.take(300).forEach { it.verify() }
|
||||
|
||||
val t1 = System.nanoTime()
|
||||
sample.forEach { it.verify() }
|
||||
val seqNanos = System.nanoTime() - t1
|
||||
println(
|
||||
" sequential (1 thread): %.1fms -> %.0f events/s (%.1fµs/event)"
|
||||
.format(seqNanos / 1e6, sample.size * 1e9 / seqNanos, seqNanos / 1e3 / sample.size),
|
||||
)
|
||||
|
||||
val cores = Runtime.getRuntime().availableProcessors()
|
||||
val parNanos =
|
||||
runBlocking(Dispatchers.Default) {
|
||||
val start = System.nanoTime()
|
||||
sample
|
||||
.chunked((sample.size + cores - 1) / cores)
|
||||
.map { chunk -> async { chunk.forEach { it.verify() } } }
|
||||
.awaitAll()
|
||||
System.nanoTime() - start
|
||||
}
|
||||
println(
|
||||
" parallel ($cores threads): %.1fms -> %.0f events/s (%.1fx speedup)"
|
||||
.format(parNanos / 1e6, sample.size * 1e9 / parNanos, seqNanos.toDouble() / parNanos),
|
||||
)
|
||||
}
|
||||
|
||||
// -----------------------------------------------------------------
|
||||
// The gated test
|
||||
// -----------------------------------------------------------------
|
||||
|
||||
@Test
|
||||
fun productionReceiverBenchmark() {
|
||||
if (System.getenv("PROD_RELAY_BENCH") == null && System.getProperty("prodRelayBench") == null) {
|
||||
println("ProductionReceiverBenchmark skipped. Set PROD_RELAY_BENCH=1 to run against live relays.")
|
||||
return
|
||||
}
|
||||
|
||||
val httpClient =
|
||||
OkHttpClient
|
||||
.Builder()
|
||||
.connectTimeout(15, TimeUnit.SECONDS)
|
||||
.readTimeout(60, TimeUnit.SECONDS)
|
||||
.pingInterval(30, TimeUnit.SECONDS)
|
||||
.build()
|
||||
|
||||
val relays = RELAYS.map { it.normalizeRelayUrl() }
|
||||
|
||||
runBlocking {
|
||||
val capturedFrames = Collections.synchronizedList(ArrayList<String>(30_000))
|
||||
|
||||
// Warmup (discarded): JIT-compile the parse/verify/dispatch path and
|
||||
// open the OkHttp pools so the first measured scenario isn't paying
|
||||
// cold-start costs that the later ones don't.
|
||||
runScenario(
|
||||
"warmup",
|
||||
httpClient,
|
||||
relays,
|
||||
relays.associateWith { listOf(Filter(kinds = listOf(1), limit = 150)) },
|
||||
VerifyMode.INLINE,
|
||||
quiet = true,
|
||||
)
|
||||
|
||||
// S1: home-feed-style firehose, current app behavior (inline verify)
|
||||
val s1 =
|
||||
runScenario(
|
||||
"kind-1 firehose, limit 500/relay",
|
||||
httpClient,
|
||||
relays,
|
||||
relays.associateWith { listOf(Filter(kinds = listOf(1), limit = 500)) },
|
||||
VerifyMode.INLINE,
|
||||
captureFrames = capturedFrames,
|
||||
)
|
||||
|
||||
// S1b: same filter, verification offloaded to a parallel pool
|
||||
runScenario(
|
||||
"kind-1 firehose, limit 500/relay",
|
||||
httpClient,
|
||||
relays,
|
||||
relays.associateWith { listOf(Filter(kinds = listOf(1), limit = 500)) },
|
||||
VerifyMode.PARALLEL,
|
||||
)
|
||||
|
||||
// S2: notifications-style filter for a busy pubkey
|
||||
runScenario(
|
||||
"notifications for busy pubkey (kinds 1,6,7,9735)",
|
||||
httpClient,
|
||||
relays,
|
||||
relays.associateWith {
|
||||
listOf(
|
||||
Filter(
|
||||
kinds = listOf(1, 6, 7, 9735),
|
||||
tags = mapOf("p" to listOf(BUSY_PUBKEY)),
|
||||
limit = 500,
|
||||
),
|
||||
)
|
||||
},
|
||||
VerifyMode.INLINE,
|
||||
)
|
||||
|
||||
// S3: metadata burst — the app-startup pattern: fetch kind 0/10002
|
||||
// for every author seen in the feed. Big bursts, tiny events.
|
||||
val authors =
|
||||
s1.capturedEvents
|
||||
.map { it.pubKey }
|
||||
.distinct()
|
||||
.take(300)
|
||||
if (authors.isNotEmpty()) {
|
||||
runScenario(
|
||||
"metadata burst (kinds 0,10002) for ${authors.size} authors",
|
||||
httpClient,
|
||||
relays,
|
||||
relays.associateWith { listOf(Filter(kinds = listOf(0, 10002), authors = authors)) },
|
||||
VerifyMode.INLINE,
|
||||
)
|
||||
runScenario(
|
||||
"metadata burst (kinds 0,10002) for ${authors.size} authors",
|
||||
httpClient,
|
||||
relays,
|
||||
relays.associateWith { listOf(Filter(kinds = listOf(0, 10002), authors = authors)) },
|
||||
VerifyMode.PARALLEL,
|
||||
)
|
||||
}
|
||||
|
||||
// Offline ceilings from the captured production data
|
||||
offlineParseBench(ArrayList(capturedFrames))
|
||||
offlineVerifyBench(s1.capturedEvents)
|
||||
}
|
||||
|
||||
httpClient.dispatcher.executorService.shutdown()
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user