mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-07-22 07:48:27 +00:00
perf(relay): remove live-path allocations in LiveEventStore + FilterIndex
The remaining SmallReqFloorBenchmark waste, on the per-row replay, the per-event live fanout, and the per-accepted-event index probe: - LiveEventStore replay dedupe: a SeenIds holder with an inline lock replaces the local-fn-plus-lambda that allocated one closure per streamed row (and again per live delivery). Its HashSet is created empty so the JVM defers the backing table to the first add — a 0-row replay no longer allocates a 1024-slot table (was ~4 MB across the benchmark's 1000 idle subs). - Live fanout serializes the event body once and passes it through onEachLive(event, body); RelaySession splices it into the per-sub frame prefix. An event matching N live subscriptions paid N identical Jackson passes before; now one. queryRaw's onEachLive signature gains the body arg (EventSourceBackend default serializes inline, no cross-sub memo, no regression). Measured: fanout 1->200 live subs 0.50 ms (2.5 us/sub). - FilterIndex holds subscribers in one persistent map per dimension, so candidatesFor (once per accepted ingest event) probes with the event's own fields and allocates no IdKey/AuthorKey/KindKey/TagKey wrappers; BucketKey now lives only in the rare register/unregister bookkeeping. SmallReqFloorBenchmark grows a fanout stage (200 live subs, one submit) to anchor the fanout number; it drives `live` directly and guards the await with withTimeout so a future fanout regression fails fast. Verified: quartz relay.server + FilterIndex suites (110 tests), SmallReqFloorBenchmark, geode suite (126 tests). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01RzkoN3SJHCZWRAiadXXG4w
This commit is contained in:
@@ -115,19 +115,25 @@ class FilterIndex<S : Any> {
|
||||
private object Unindexed : BucketKey
|
||||
|
||||
/**
|
||||
* Single immutable snapshot. [buckets] maps a key to the set of
|
||||
* subscribers registered under it; [assignments] is the reverse
|
||||
* map used by [unregister] to find a subscriber's keys without
|
||||
* scanning every bucket.
|
||||
* Single immutable snapshot. Subscribers are held in one map per
|
||||
* indexable dimension so [candidatesFor] — called once per accepted
|
||||
* ingest event, the hot read — can probe each dimension with the
|
||||
* event's own field (`event.id`, `event.pubKey`, `tag[0]`/`tag[1]`,
|
||||
* `event.kind`) and allocate no key-wrapper objects. [assignments] is
|
||||
* the reverse map ([S] → the [BucketKey]s it occupies) used by
|
||||
* [unregister]; the wrappers live only here, built on the rare
|
||||
* register path.
|
||||
*
|
||||
* Persistent (HAMT) maps/sets: a register/unregister produces the
|
||||
* next snapshot in O(keys × log S) with structural sharing, instead
|
||||
* of copying both full maps — registration happens on every REQ
|
||||
* open/close, so with S live subscriptions the full copy was
|
||||
* O(S) work and O(S) allocation per REQ.
|
||||
* of copying full maps — registration happens on every REQ open/close.
|
||||
*/
|
||||
private data class State<S>(
|
||||
val buckets: PersistentMap<BucketKey, PersistentSet<S>> = persistentHashMapOf(),
|
||||
val ids: PersistentMap<HexKey, PersistentSet<S>> = persistentHashMapOf(),
|
||||
val authors: PersistentMap<HexKey, PersistentSet<S>> = persistentHashMapOf(),
|
||||
val tags: PersistentMap<String, PersistentMap<String, PersistentSet<S>>> = persistentHashMapOf(),
|
||||
val kinds: PersistentMap<Int, PersistentSet<S>> = persistentHashMapOf(),
|
||||
val unindexed: PersistentSet<S> = persistentHashSetOf(),
|
||||
val assignments: PersistentMap<S, PersistentSet<BucketKey>> = persistentHashMapOf(),
|
||||
)
|
||||
|
||||
@@ -193,19 +199,23 @@ class FilterIndex<S : Any> {
|
||||
while (true) {
|
||||
val current = state.load()
|
||||
val keys = current.assignments[subscriber] ?: return
|
||||
var newBuckets = current.buckets
|
||||
var ids = current.ids
|
||||
var authors = current.authors
|
||||
var tags = current.tags
|
||||
var kinds = current.kinds
|
||||
var unindexed = current.unindexed
|
||||
for (key in keys) {
|
||||
val cur = newBuckets[key] ?: continue
|
||||
val next = cur.remove(subscriber)
|
||||
newBuckets =
|
||||
if (next.isEmpty()) {
|
||||
newBuckets.remove(key)
|
||||
} else {
|
||||
newBuckets.put(key, next)
|
||||
}
|
||||
when (key) {
|
||||
is IdKey -> ids = ids.removeSub(key.id, subscriber)
|
||||
is AuthorKey -> authors = authors.removeSub(key.author, subscriber)
|
||||
is KindKey -> kinds = kinds.removeSub(key.kind, subscriber)
|
||||
is TagKey -> tags = tags.removeTagSub(key.letter, key.value, subscriber)
|
||||
Unindexed -> unindexed = unindexed.remove(subscriber)
|
||||
}
|
||||
}
|
||||
val newAssignments = current.assignments.remove(subscriber)
|
||||
if (state.compareAndSet(current, State(newBuckets, newAssignments))) return
|
||||
val next =
|
||||
State(ids, authors, tags, kinds, unindexed, current.assignments.remove(subscriber))
|
||||
if (state.compareAndSet(current, next)) return
|
||||
}
|
||||
}
|
||||
|
||||
@@ -215,19 +225,22 @@ class FilterIndex<S : Any> {
|
||||
* candidate to handle negative constraints.
|
||||
*
|
||||
* Iteration order is insertion-stable per call but otherwise
|
||||
* unspecified.
|
||||
* unspecified. Allocates only the result set — dimensions are
|
||||
* probed with the event's own fields, no key wrappers.
|
||||
*/
|
||||
fun candidatesFor(event: Event): Set<S> {
|
||||
val s = state.load()
|
||||
if (s.buckets.isEmpty()) return emptySet()
|
||||
if (s.assignments.isEmpty()) return emptySet()
|
||||
val result = LinkedHashSet<S>()
|
||||
s.buckets[Unindexed]?.let { result.addAll(it) }
|
||||
s.buckets[IdKey(event.id)]?.let { result.addAll(it) }
|
||||
s.buckets[AuthorKey(event.pubKey)]?.let { result.addAll(it) }
|
||||
s.buckets[KindKey(event.kind)]?.let { result.addAll(it) }
|
||||
for (tag in event.tags) {
|
||||
if (tag.size >= 2 && tag[0].length == 1) {
|
||||
s.buckets[TagKey(tag[0], tag[1])]?.let { result.addAll(it) }
|
||||
if (s.unindexed.isNotEmpty()) result.addAll(s.unindexed)
|
||||
s.ids[event.id]?.let { result.addAll(it) }
|
||||
s.authors[event.pubKey]?.let { result.addAll(it) }
|
||||
s.kinds[event.kind]?.let { result.addAll(it) }
|
||||
if (s.tags.isNotEmpty()) {
|
||||
for (tag in event.tags) {
|
||||
if (tag.size >= 2 && tag[0].length == 1) {
|
||||
s.tags[tag[0]]?.get(tag[1])?.let { result.addAll(it) }
|
||||
}
|
||||
}
|
||||
}
|
||||
return result
|
||||
@@ -251,16 +264,72 @@ class FilterIndex<S : Any> {
|
||||
val keySet = keys.toPersistentHashSet()
|
||||
while (true) {
|
||||
val current = state.load()
|
||||
var newBuckets = current.buckets
|
||||
var ids = current.ids
|
||||
var authors = current.authors
|
||||
var tags = current.tags
|
||||
var kinds = current.kinds
|
||||
var unindexed = current.unindexed
|
||||
for (key in keySet) {
|
||||
val cur = newBuckets[key] ?: persistentHashSetOf()
|
||||
val next = cur.add(subscriber)
|
||||
if (next !== cur) newBuckets = newBuckets.put(key, next)
|
||||
when (key) {
|
||||
is IdKey -> ids = ids.addSub(key.id, subscriber)
|
||||
is AuthorKey -> authors = authors.addSub(key.author, subscriber)
|
||||
is KindKey -> kinds = kinds.addSub(key.kind, subscriber)
|
||||
is TagKey -> tags = tags.addTagSub(key.letter, key.value, subscriber)
|
||||
Unindexed -> unindexed = unindexed.add(subscriber)
|
||||
}
|
||||
}
|
||||
val existing = current.assignments[subscriber]
|
||||
val merged = existing?.addAll(keySet) ?: keySet
|
||||
val newAssignments = current.assignments.put(subscriber, merged)
|
||||
if (state.compareAndSet(current, State(newBuckets, newAssignments))) return
|
||||
val next = State(ids, authors, tags, kinds, unindexed, current.assignments.put(subscriber, merged))
|
||||
if (state.compareAndSet(current, next)) return
|
||||
}
|
||||
}
|
||||
|
||||
// Per-dimension add/remove of one subscriber, returning the same map
|
||||
// instance when nothing changed so the CAS builds minimal new nodes.
|
||||
private fun <K> PersistentMap<K, PersistentSet<S>>.addSub(
|
||||
key: K,
|
||||
sub: S,
|
||||
): PersistentMap<K, PersistentSet<S>> {
|
||||
val cur = this[key] ?: persistentHashSetOf()
|
||||
val next = cur.add(sub)
|
||||
return if (next === cur) this else put(key, next)
|
||||
}
|
||||
|
||||
private fun <K> PersistentMap<K, PersistentSet<S>>.removeSub(
|
||||
key: K,
|
||||
sub: S,
|
||||
): PersistentMap<K, PersistentSet<S>> {
|
||||
val cur = this[key] ?: return this
|
||||
val next = cur.remove(sub)
|
||||
return when {
|
||||
next === cur -> this
|
||||
next.isEmpty() -> remove(key)
|
||||
else -> put(key, next)
|
||||
}
|
||||
}
|
||||
|
||||
private fun PersistentMap<String, PersistentMap<String, PersistentSet<S>>>.addTagSub(
|
||||
letter: String,
|
||||
value: String,
|
||||
sub: S,
|
||||
): PersistentMap<String, PersistentMap<String, PersistentSet<S>>> {
|
||||
val inner = this[letter] ?: persistentHashMapOf()
|
||||
val newInner = inner.addSub(value, sub)
|
||||
return if (newInner === inner) this else put(letter, newInner)
|
||||
}
|
||||
|
||||
private fun PersistentMap<String, PersistentMap<String, PersistentSet<S>>>.removeTagSub(
|
||||
letter: String,
|
||||
value: String,
|
||||
sub: S,
|
||||
): PersistentMap<String, PersistentMap<String, PersistentSet<S>>> {
|
||||
val inner = this[letter] ?: return this
|
||||
val newInner = inner.removeSub(value, sub)
|
||||
return when {
|
||||
newInner === inner -> this
|
||||
newInner.isEmpty() -> remove(letter)
|
||||
else -> put(letter, newInner)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -340,7 +340,20 @@ class RelaySession(
|
||||
},
|
||||
)
|
||||
},
|
||||
onEachLive = { event -> send(EventMessage(cmd.subId, event)) },
|
||||
// Live events arrive with their wire body already
|
||||
// serialized (once per event, shared across every
|
||||
// matching subscription): splice it into the same
|
||||
// per-sub frame prefix as the stored replay, no
|
||||
// per-event EventMessage or re-serialize.
|
||||
onEachLive = { _, body ->
|
||||
sendRaw(
|
||||
buildString(framePrefix.length + body.length + 1) {
|
||||
append(framePrefix)
|
||||
append(body)
|
||||
append(']')
|
||||
},
|
||||
)
|
||||
},
|
||||
onEose = { send(EoseMessage(cmd.subId)) },
|
||||
)
|
||||
}
|
||||
|
||||
@@ -74,13 +74,57 @@ class LiveEventStore(
|
||||
* One live REQ subscription. Carries the filters (for the
|
||||
* post-index `match` re-check needed for negative constraints
|
||||
* like `since` / `until` / `tagsAll`) and the delivery callback
|
||||
* the index dispatches into. Identity-keyed inside [FilterIndex].
|
||||
* the index dispatches into. [deliver] receives the event and its
|
||||
* pre-serialized wire body (memoized once per fanout across all
|
||||
* matching subscribers). Identity-keyed inside [FilterIndex].
|
||||
*/
|
||||
private class LiveSubscription(
|
||||
val filters: List<Filter>,
|
||||
val deliver: (Event) -> Unit,
|
||||
val deliver: (Event, String) -> Unit,
|
||||
)
|
||||
|
||||
/**
|
||||
* Replay-dedupe set for one REQ. During the historical replay the
|
||||
* store's ids are [record]ed here so the concurrent live path can
|
||||
* drop an event the replay also emitted; after EOSE the set is
|
||||
* [release]d and the live path forwards everything.
|
||||
*
|
||||
* Written from the replay coroutine and read from the [IngestQueue]
|
||||
* drain coroutine (via `fanout`), so every access takes a tiny spin
|
||||
* lock — [locked] is `inline`, so the per-row `record` / `isDuplicate`
|
||||
* calls allocate no closure. The backing `HashSet` is created empty
|
||||
* up front (so the register-before-replay race guarantee holds) but
|
||||
* the JVM defers its table allocation to the first `add`, so a
|
||||
* zero-row replay costs only the empty set object, not a sized table.
|
||||
* It MUST stay a mutable set under a lock, never a copy-on-add
|
||||
* immutable set — `set + id` per row made large replays O(n²).
|
||||
*/
|
||||
private class SeenIds {
|
||||
private val lock = AtomicBoolean(false)
|
||||
private var ids: HashSet<String>? = HashSet()
|
||||
|
||||
private inline fun <R> locked(block: () -> R): R {
|
||||
while (lock.exchange(true)) {
|
||||
while (lock.load()) { }
|
||||
}
|
||||
try {
|
||||
return block()
|
||||
} finally {
|
||||
lock.store(false)
|
||||
}
|
||||
}
|
||||
|
||||
fun record(id: String) {
|
||||
locked { ids?.add(id) }
|
||||
}
|
||||
|
||||
fun isDuplicate(id: String): Boolean = locked { ids?.contains(id) ?: false }
|
||||
|
||||
fun release() {
|
||||
locked { ids = null }
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Fire-and-forget enqueue: hand [event] to the [IngestQueue] and
|
||||
* fire [onComplete] once the writer's batch has a per-row
|
||||
@@ -149,9 +193,17 @@ class LiveEventStore(
|
||||
* batch writer.
|
||||
*/
|
||||
private fun fanout(event: Event) {
|
||||
for (sub in index.candidatesFor(event)) {
|
||||
val candidates = index.candidatesFor(event)
|
||||
if (candidates.isEmpty()) return
|
||||
// Serialize the wire body at most once for this event, no matter
|
||||
// how many subscriptions match it — the old path re-serialized the
|
||||
// whole event per matching subscriber, so a note landing in N live
|
||||
// feeds paid N identical Jackson passes. Lazy so a fanout that
|
||||
// matches nothing (index over-approximates) serializes nothing.
|
||||
var body: String? = null
|
||||
for (sub in candidates) {
|
||||
if (sub.filters.any { it.match(event) }) {
|
||||
sub.deliver(event)
|
||||
sub.deliver(event, body ?: event.toJson().also { body = it })
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -178,61 +230,33 @@ class LiveEventStore(
|
||||
onEose: () -> Unit,
|
||||
) {
|
||||
drainFtsIfSearching(filters)
|
||||
// During the historical replay, record ids the store has
|
||||
// emitted so the live path can dedupe. The index registers
|
||||
// *before* the replay starts (otherwise an event accepted
|
||||
// mid-replay would slip past the live path entirely — same
|
||||
// race the previous SharedFlow-based implementation closed
|
||||
// with `onSubscription`).
|
||||
//
|
||||
// The set is read from the [IngestQueue] drain coroutine (in
|
||||
// `deliver`, called synchronously from `fanout`) and written
|
||||
// from this coroutine (the historical-replay closure below),
|
||||
// so access is guarded by a tiny spin lock (contains/add,
|
||||
// never I/O). It MUST be a mutable set under a lock, not an
|
||||
// immutable Set under an AtomicReference with copy-on-add:
|
||||
// `set + id` copies the whole set per streamed event, which
|
||||
// made large replays accidentally O(n²) — a 100k-event REQ
|
||||
// crawled at ~700 events/s and the rate degraded as the
|
||||
// response grew (see the plan doc's giant-REQ finding).
|
||||
//
|
||||
// Once cleared to null after EOSE, `deliver` short-circuits
|
||||
// and every live event is forwarded.
|
||||
val seenLock = AtomicBoolean(false)
|
||||
var seenIds: HashSet<String>? = HashSet(1024)
|
||||
|
||||
fun <R> seenLocked(block: () -> R): R {
|
||||
while (seenLock.exchange(true)) {
|
||||
while (seenLock.load()) { }
|
||||
}
|
||||
try {
|
||||
return block()
|
||||
} finally {
|
||||
seenLock.store(false)
|
||||
}
|
||||
}
|
||||
// The index registers *before* the replay starts (otherwise an
|
||||
// event accepted mid-replay would slip past the live path entirely
|
||||
// — same race the previous SharedFlow-based implementation closed
|
||||
// with `onSubscription`), and [SeenIds] bridges the two coroutines:
|
||||
// the replay records ids here, the live `deliver` drops duplicates,
|
||||
// and after EOSE the set is released so every live event forwards.
|
||||
val seen = SeenIds()
|
||||
|
||||
val sub =
|
||||
LiveSubscription(
|
||||
filters = filters,
|
||||
deliver = { event ->
|
||||
val duplicate = seenLocked { seenIds?.contains(event.id) ?: false }
|
||||
if (duplicate) return@LiveSubscription
|
||||
onEach(event)
|
||||
deliver = { event, _ ->
|
||||
if (!seen.isDuplicate(event.id)) onEach(event)
|
||||
},
|
||||
)
|
||||
|
||||
index.register(filters, sub)
|
||||
try {
|
||||
store.query<Event>(filters.strippingSearchExtensions()) { event ->
|
||||
seenLocked { seenIds?.add(event.id) }
|
||||
seen.record(event.id)
|
||||
onEach(event)
|
||||
}
|
||||
onEose()
|
||||
// Drop the dedupe set so the live path stops paying for
|
||||
// it. From this point the index drives delivery and
|
||||
// duplicates are no longer possible.
|
||||
seenLocked { seenIds = null }
|
||||
seen.release()
|
||||
// Suspend until the caller's coroutine is cancelled
|
||||
// (e.g. NIP-01 CLOSE or connection drop). The `finally`
|
||||
// unregisters from the index.
|
||||
@@ -255,42 +279,28 @@ class LiveEventStore(
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
onEachStored: (RawEvent) -> Unit,
|
||||
onEachLive: (Event) -> Unit,
|
||||
onEachLive: (Event, String) -> Unit,
|
||||
onEose: () -> Unit,
|
||||
) {
|
||||
drainFtsIfSearching(filters)
|
||||
val seenLock = AtomicBoolean(false)
|
||||
var seenIds: HashSet<String>? = HashSet(1024)
|
||||
|
||||
fun <R> seenLocked(block: () -> R): R {
|
||||
while (seenLock.exchange(true)) {
|
||||
while (seenLock.load()) { }
|
||||
}
|
||||
try {
|
||||
return block()
|
||||
} finally {
|
||||
seenLock.store(false)
|
||||
}
|
||||
}
|
||||
val seen = SeenIds()
|
||||
|
||||
val sub =
|
||||
LiveSubscription(
|
||||
filters = filters,
|
||||
deliver = { event ->
|
||||
val duplicate = seenLocked { seenIds?.contains(event.id) ?: false }
|
||||
if (duplicate) return@LiveSubscription
|
||||
onEachLive(event)
|
||||
deliver = { event, body ->
|
||||
if (!seen.isDuplicate(event.id)) onEachLive(event, body)
|
||||
},
|
||||
)
|
||||
|
||||
index.register(filters, sub)
|
||||
try {
|
||||
store.rawQuery(filters.strippingSearchExtensions()) { raw ->
|
||||
seenLocked { seenIds?.add(raw.id) }
|
||||
seen.record(raw.id)
|
||||
onEachStored(raw)
|
||||
}
|
||||
onEose()
|
||||
seenLocked { seenIds = null }
|
||||
seen.release()
|
||||
awaitCancellation()
|
||||
} finally {
|
||||
index.unregister(sub)
|
||||
|
||||
@@ -78,9 +78,9 @@ interface SessionBackend {
|
||||
ctx: RequestContext,
|
||||
filters: List<Filter>,
|
||||
onEachStored: (RawEvent) -> Unit,
|
||||
onEachLive: (Event) -> Unit,
|
||||
onEachLive: (Event, String) -> Unit,
|
||||
onEose: () -> Unit,
|
||||
): Unit = query(ctx, filters, onEachLive, onEose)
|
||||
): Unit = query(ctx, filters, { onEachLive(it, it.toJson()) }, onEose)
|
||||
|
||||
/** Answers a NIP-45 COUNT with an exact cardinality for the caller in [ctx]. */
|
||||
suspend fun count(
|
||||
|
||||
@@ -38,6 +38,8 @@ import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.withTimeout
|
||||
import java.util.concurrent.atomic.AtomicInteger
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
|
||||
@@ -67,6 +69,7 @@ class SmallReqFloorBenchmark {
|
||||
const val ROUNDS = 400
|
||||
const val WARMUP = 100
|
||||
const val IDLE_SUBS = 1_000
|
||||
const val FANOUT_SUBS = 200
|
||||
}
|
||||
|
||||
private fun hexId(seed: Int): String = seed.toString(16).padStart(64, '0')
|
||||
@@ -136,7 +139,7 @@ class SmallReqFloorBenchmark {
|
||||
ctx = ctx,
|
||||
filters = listOf(filterFor(round)),
|
||||
onEachStored = {},
|
||||
onEachLive = {},
|
||||
onEachLive = { _, _ -> },
|
||||
onEose = { eose.complete(System.nanoTime() - t0) },
|
||||
)
|
||||
}
|
||||
@@ -163,7 +166,7 @@ class SmallReqFloorBenchmark {
|
||||
ctx = ctx,
|
||||
filters = listOf(Filter(authors = listOf(hexId(1_000_000 + i)), kinds = listOf(1), limit = 1)),
|
||||
onEachStored = {},
|
||||
onEachLive = {},
|
||||
onEachLive = { _, _ -> },
|
||||
onEose = { ready.complete(Unit) },
|
||||
)
|
||||
}
|
||||
@@ -192,6 +195,54 @@ class SmallReqFloorBenchmark {
|
||||
val c = LongArray(ROUNDS)
|
||||
repeat(ROUNDS) { c[it] = timeSession(it) }
|
||||
|
||||
// --- fanout: one live event → FANOUT_SUBS live subscriptions ---
|
||||
// All subs register on `live` directly (via queryRaw, same backend
|
||||
// we submit into) and filter an author with no stored events (0-row
|
||||
// replay, then park live). Submitting one matching event fans out to
|
||||
// every sub; the body is serialized once and spliced per sub, so
|
||||
// this measures the shared-serialization path (#2). skipVerify so
|
||||
// the synthetic sig is accepted. Fewer rounds than A–C: each round
|
||||
// is FANOUT_SUBS deliveries and a real group-commit insert.
|
||||
val fanAuthor = hexId(9_000_001)
|
||||
val delivered = AtomicInteger(0)
|
||||
var fanDone = CompletableDeferred<Long>()
|
||||
var fanStart = 0L
|
||||
val fanJobs =
|
||||
(0 until FANOUT_SUBS).map {
|
||||
val ready = CompletableDeferred<Unit>()
|
||||
val job =
|
||||
scope.launch(start = CoroutineStart.UNDISPATCHED) {
|
||||
live.queryRaw(
|
||||
ctx = ctx,
|
||||
filters = listOf(Filter(authors = listOf(fanAuthor), kinds = listOf(1))),
|
||||
onEachStored = {},
|
||||
onEachLive = { _, _ ->
|
||||
if (delivered.incrementAndGet() == FANOUT_SUBS) {
|
||||
fanDone.complete(System.nanoTime() - fanStart)
|
||||
}
|
||||
},
|
||||
onEose = { ready.complete(Unit) },
|
||||
)
|
||||
}
|
||||
ready.await()
|
||||
job
|
||||
}
|
||||
val fanRounds = 60
|
||||
val fanWarmup = 15
|
||||
val fan = LongArray(fanRounds)
|
||||
var fanSeq = 0
|
||||
repeat(fanWarmup + fanRounds) { r ->
|
||||
delivered.set(0)
|
||||
fanDone = CompletableDeferred()
|
||||
val ev = EventFactory.create<Event>(hexId(9_500_000 + fanSeq), fanAuthor, 1_700_000_000L + fanSeq, 1, emptyArray(), "fanout $fanSeq", sig)
|
||||
fanSeq++
|
||||
fanStart = System.nanoTime()
|
||||
live.submit(ev, skipVerify = true) {}
|
||||
val nanos = withTimeout(30_000) { fanDone.await() }
|
||||
if (r >= fanWarmup) fan[r - fanWarmup] = nanos
|
||||
}
|
||||
fanJobs.forEach { it.cancel() }
|
||||
|
||||
assertEquals(true, rowsA > 0, "author filters must return rows")
|
||||
|
||||
val mA = median(a)
|
||||
@@ -203,6 +254,8 @@ class SmallReqFloorBenchmark {
|
||||
println(" B backend queryRaw→EOSE: ${"%6.3f".format(mB)} ms (live machinery +${"%6.3f".format(mB - mA)})")
|
||||
println(" B@${IDLE_SUBS} idle subs: ${"%6.3f".format(mB1k)} ms (population cost +${"%6.3f".format(mB1k - mB)})")
|
||||
println(" C session REQ→EOSE: ${"%6.3f".format(mC)} ms (dispatch+frames +${"%6.3f".format(mC - mB)})")
|
||||
val mFan = median(fan)
|
||||
println(" fanout 1→$FANOUT_SUBS live subs: ${"%6.3f".format(mFan)} ms (${"%.2f".format(mFan * 1000 / FANOUT_SUBS)} µs/sub; body serialized once)")
|
||||
|
||||
server.close()
|
||||
scope.cancel()
|
||||
|
||||
Reference in New Issue
Block a user