mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 19:53:08 +00:00
test(perf): add ingest-latency and NEG-MSG serialization benchmarks
Measurement harnesses for the two remaining relayBench gaps, isolating each cost so a fix can be judged on the delta: - IngestLatencyBenchmark: times single-event submit→onComplete through the group-commit IngestQueue vs a direct batchInsert. The delta is the pipeline's coroutine-handoff overhead — the receipt→queryable latency the 1M run measured geode losing (4.68ms vs strfry 2.32ms). - NegMsgSerializationBenchmark: times the NEG-MSG wire path (Hex.encode → MessageKSerializer JsonElement tree → UTF-8) against a direct StringBuilder build, asserting byte-identical output. Isolates the ~40% serialization slice of the server reconcile the JFR flagged. Both compile and are self-contained; results + fixes to follow. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_012EZeWww5TJnzBZKPoc6mvU
This commit is contained in:
+128
@@ -0,0 +1,128 @@
|
||||
/*
|
||||
* 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.relay.server.backend.IngestQueue
|
||||
import com.vitorpamplona.quartz.nip01Core.store.IEventStore
|
||||
import com.vitorpamplona.quartz.nip01Core.store.sqlite.DefaultIndexingStrategy
|
||||
import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
|
||||
import com.vitorpamplona.quartz.utils.EventFactory
|
||||
import kotlinx.coroutines.CompletableDeferred
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlin.test.Test
|
||||
|
||||
/**
|
||||
* Isolates the receipt➜queryable latency that relayBench measured geode
|
||||
* losing to strfry (1M corpus: geode 4.68 ms p50 vs strfry 2.32 ms). The
|
||||
* probe there publishes ONE event to an idle relay and polls REQ until it
|
||||
* comes back, so what it times is a single event's trip through the ingest
|
||||
* pipeline — where geode's group-commit [IngestQueue] pays two channel
|
||||
* handoffs (submit→verifier→writer) for zero batching benefit.
|
||||
*
|
||||
* This times `submit → onComplete` (which fires *after* COMMIT, so it's the
|
||||
* moment the row is queryable) for isolated single events, against a direct
|
||||
* `store.batchInsert(listOf(e))` — the delta is exactly the pipeline's
|
||||
* coroutine-handoff overhead, the thing a single-event fast path would save.
|
||||
*
|
||||
* Prints percentiles; no speed assertion (container-noisy).
|
||||
*/
|
||||
class IngestLatencyBenchmark {
|
||||
private val hexChars = "0123456789abcdef".toCharArray()
|
||||
|
||||
private fun mix(seed: Long): Long {
|
||||
var z = seed + -0x61c8864680b583ebL
|
||||
z = (z xor (z ushr 30)) * -0x40a7b892e31b1a47L
|
||||
z = (z xor (z ushr 27)) * -0x6b2fb644ecceee15L
|
||||
return z xor (z ushr 31)
|
||||
}
|
||||
|
||||
private fun idFor(index: Int): String {
|
||||
val out = CharArray(64)
|
||||
for (w in 0 until 4) {
|
||||
val v = mix(index.toLong() * 4 + w)
|
||||
for (b in 0 until 8) {
|
||||
val byte = ((v ushr (b * 8)) and 0xFF).toInt()
|
||||
out[(w * 8 + b) * 2] = hexChars[byte ushr 4]
|
||||
out[(w * 8 + b) * 2 + 1] = hexChars[byte and 0xF]
|
||||
}
|
||||
}
|
||||
return String(out)
|
||||
}
|
||||
|
||||
private val pubkey = "00".repeat(32)
|
||||
private val sig = "0".repeat(128)
|
||||
|
||||
private fun event(i: Int): Event = EventFactory.create(idFor(i), pubkey, 1_700_000_000L + i, 1, emptyArray(), "note $i", sig)
|
||||
|
||||
private fun pct(
|
||||
sorted: DoubleArray,
|
||||
p: Double,
|
||||
) = sorted[(p * (sorted.size - 1)).toInt()]
|
||||
|
||||
@Test
|
||||
fun singleEventVisibilityLatency() =
|
||||
runBlocking {
|
||||
val warmup = 500
|
||||
val runs = 3_000
|
||||
|
||||
// ── Path A: through the group-commit IngestQueue (production). ──
|
||||
val storeA = EventStore(dbName = null, indexStrategy = DefaultIndexingStrategy())
|
||||
val queue = IngestQueue(storeA, Dispatchers.Default + SupervisorJob())
|
||||
val queueLat = DoubleArray(runs)
|
||||
var idx = 0
|
||||
repeat(warmup + runs) { i ->
|
||||
val done = CompletableDeferred<IEventStore.InsertOutcome>()
|
||||
val t0 = System.nanoTime()
|
||||
queue.submit(event(idx++)) { done.complete(it) }
|
||||
done.await()
|
||||
val dt = (System.nanoTime() - t0) / 1e6
|
||||
if (i >= warmup) queueLat[i - warmup] = dt
|
||||
}
|
||||
queue.close()
|
||||
storeA.close()
|
||||
|
||||
// ── Path B: direct single-event batchInsert (no pipeline). ──
|
||||
val storeB = EventStore(dbName = null, indexStrategy = DefaultIndexingStrategy())
|
||||
val directLat = DoubleArray(runs)
|
||||
repeat(warmup + runs) { i ->
|
||||
val t0 = System.nanoTime()
|
||||
storeB.batchInsert(listOf(event(idx++)))
|
||||
val dt = (System.nanoTime() - t0) / 1e6
|
||||
if (i >= warmup) directLat[i - warmup] = dt
|
||||
}
|
||||
storeB.close()
|
||||
|
||||
queueLat.sort()
|
||||
directLat.sort()
|
||||
println("─ IngestLatencyBenchmark: single isolated events (idle relay), $runs samples ─")
|
||||
println(" path p50 p90 p99")
|
||||
println(
|
||||
" IngestQueue (submit→OK) ${"%6.3f".format(pct(queueLat, 0.50))} ms ${"%6.3f".format(pct(queueLat, 0.90))} ms ${"%6.3f".format(pct(queueLat, 0.99))} ms",
|
||||
)
|
||||
println(
|
||||
" direct batchInsert ${"%6.3f".format(pct(directLat, 0.50))} ms ${"%6.3f".format(pct(directLat, 0.90))} ms ${"%6.3f".format(pct(directLat, 0.99))} ms",
|
||||
)
|
||||
println(" pipeline overhead p50: ${"%.3f".format(pct(queueLat, 0.50) - pct(directLat, 0.50))} ms")
|
||||
}
|
||||
}
|
||||
+121
@@ -0,0 +1,121 @@
|
||||
/*
|
||||
* 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.nip77Negentropy.NegMsgMessage
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
|
||||
/**
|
||||
* Isolates the per-round NEG-MSG serialization tax that geode's server-path
|
||||
* JFR attributed ~40% of reconcile time to: a reconcile produces a binary
|
||||
* frame, `Hex.encode`s it to a ~1 MB hex String, then the outgoing [Message]
|
||||
* is turned into wire JSON by [com.vitorpamplona.quartz.nip01Core.kotlinSerialization.MessageKSerializer],
|
||||
* which builds a `JsonElement` tree (wrapping the giant hex string in a
|
||||
* `JsonPrimitive`) and re-serializes it — scanning every hex char for JSON
|
||||
* escapes that a `[0-9a-f]` payload can never contain — before Ktor UTF-8
|
||||
* encodes it to the socket.
|
||||
*
|
||||
* The wire is trivially `["NEG-MSG","<sub>","<hex>"]`, so a direct
|
||||
* `StringBuilder` skips the tree. This measures the current `toJson()` path
|
||||
* against that direct build (both followed by the UTF-8 encode a websocket
|
||||
* text frame pays), at realistic NEG-MSG sizes, and asserts they produce
|
||||
* byte-identical wire output.
|
||||
*/
|
||||
class NegMsgSerializationBenchmark {
|
||||
/** Direct wire build — `["NEG-MSG","<sub>","<hex>"]`, hex needs no escaping. */
|
||||
private fun directWire(
|
||||
subId: String,
|
||||
hex: String,
|
||||
): String =
|
||||
buildString(hex.length + subId.length + 16) {
|
||||
append("[\"NEG-MSG\",")
|
||||
append(jsonString(subId)) // client-chosen subId: escape defensively
|
||||
append(',')
|
||||
append('"')
|
||||
append(hex) // pure hex, no escaping possible
|
||||
append("\"]")
|
||||
}
|
||||
|
||||
/** Minimal JSON string encoder for the (short, usually-safe) subId. */
|
||||
private fun jsonString(s: String): String =
|
||||
buildString(s.length + 2) {
|
||||
append('"')
|
||||
for (c in s) {
|
||||
when (c) {
|
||||
'"' -> append("\\\"")
|
||||
'\\' -> append("\\\\")
|
||||
'\n' -> append("\\n")
|
||||
'\r' -> append("\\r")
|
||||
'\t' -> append("\\t")
|
||||
else -> if (c < ' ') append("\\u%04x".format(c.code)) else append(c)
|
||||
}
|
||||
}
|
||||
append('"')
|
||||
}
|
||||
|
||||
private fun hexPayload(rawBytes: Int): String {
|
||||
val chars = "0123456789abcdef"
|
||||
val sb = StringBuilder(rawBytes * 2)
|
||||
var x = 0x9E3779B9.toInt()
|
||||
repeat(rawBytes * 2) {
|
||||
x = x * 1_103_515_245 + 12_345
|
||||
sb.append(chars[(x ushr 16) and 0xF])
|
||||
}
|
||||
return sb.toString()
|
||||
}
|
||||
|
||||
private fun bench(
|
||||
label: String,
|
||||
rawFrameBytes: Int,
|
||||
) {
|
||||
val subId = "bench-sync"
|
||||
val hex = hexPayload(rawFrameBytes)
|
||||
val msg = NegMsgMessage(subId, hex)
|
||||
|
||||
// Correctness: direct build must equal the serializer's wire output.
|
||||
assertEquals(msg.toJson(), directWire(subId, hex), "wire mismatch for $label")
|
||||
|
||||
val warmup = 200
|
||||
val runs = 500
|
||||
var sink = 0
|
||||
|
||||
repeat(warmup) { sink = sink xor msg.toJson().length xor directWire(subId, hex).length }
|
||||
|
||||
val t0 = System.nanoTime()
|
||||
repeat(runs) { sink = sink xor msg.toJson().encodeToByteArray().size }
|
||||
val curMs = (System.nanoTime() - t0) / 1e6 / runs
|
||||
|
||||
val t1 = System.nanoTime()
|
||||
repeat(runs) { sink = sink xor directWire(subId, hex).encodeToByteArray().size }
|
||||
val fastMs = (System.nanoTime() - t1) / 1e6 / runs
|
||||
|
||||
println(" %-18s cur %6.3f ms direct %6.3f ms %.1f× (sink=%d)".format(label, curMs, fastMs, curMs / fastMs, sink and 1))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun negMsgWireSerialization() {
|
||||
println("─ NegMsgSerializationBenchmark: toJson()+utf8 vs direct build+utf8 ─")
|
||||
bench("64 KiB frame", 64 * 1024)
|
||||
bench("250 KiB frame", 250 * 1024)
|
||||
bench("500 KiB frame", 500 * 1024) // strfry's frame cap
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user