From 078758888a5829c09981a102e1b372114a5f3a8e Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 21 Jul 2026 19:23:38 +0000 Subject: [PATCH] perf(relay): remove live-path allocations in LiveEventStore + FilterIndex MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 Claude-Session: https://claude.ai/code/session_01RzkoN3SJHCZWRAiadXXG4w --- .../nip01Core/relay/filters/FilterIndex.kt | 137 +++++++++++++----- .../nip01Core/relay/server/RelaySession.kt | 15 +- .../relay/server/backend/LiveEventStore.kt | 136 +++++++++-------- .../relay/server/backend/SessionBackend.kt | 4 +- .../relay/prodbench/SmallReqFloorBenchmark.kt | 57 +++++++- 5 files changed, 247 insertions(+), 102 deletions(-) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/filters/FilterIndex.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/filters/FilterIndex.kt index a5c058ec6d..bee08fc159 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/filters/FilterIndex.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/filters/FilterIndex.kt @@ -115,19 +115,25 @@ class FilterIndex { 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( - val buckets: PersistentMap> = persistentHashMapOf(), + val ids: PersistentMap> = persistentHashMapOf(), + val authors: PersistentMap> = persistentHashMapOf(), + val tags: PersistentMap>> = persistentHashMapOf(), + val kinds: PersistentMap> = persistentHashMapOf(), + val unindexed: PersistentSet = persistentHashSetOf(), val assignments: PersistentMap> = persistentHashMapOf(), ) @@ -193,19 +199,23 @@ class FilterIndex { 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 { * 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 { val s = state.load() - if (s.buckets.isEmpty()) return emptySet() + if (s.assignments.isEmpty()) return emptySet() val result = LinkedHashSet() - 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 { 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 PersistentMap>.addSub( + key: K, + sub: S, + ): PersistentMap> { + val cur = this[key] ?: persistentHashSetOf() + val next = cur.add(sub) + return if (next === cur) this else put(key, next) + } + + private fun PersistentMap>.removeSub( + key: K, + sub: S, + ): PersistentMap> { + 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>>.addTagSub( + letter: String, + value: String, + sub: S, + ): PersistentMap>> { + val inner = this[letter] ?: persistentHashMapOf() + val newInner = inner.addSub(value, sub) + return if (newInner === inner) this else put(letter, newInner) + } + + private fun PersistentMap>>.removeTagSub( + letter: String, + value: String, + sub: S, + ): PersistentMap>> { + 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) } } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelaySession.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelaySession.kt index f73493f19c..7aa1e208d6 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelaySession.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelaySession.kt @@ -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)) }, ) } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStore.kt index 1d35362134..9f5f8c4ecc 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/LiveEventStore.kt @@ -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, - 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? = HashSet() + + private inline fun 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? = HashSet(1024) - - fun 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(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, onEachStored: (RawEvent) -> Unit, - onEachLive: (Event) -> Unit, + onEachLive: (Event, String) -> Unit, onEose: () -> Unit, ) { drainFtsIfSearching(filters) - val seenLock = AtomicBoolean(false) - var seenIds: HashSet? = HashSet(1024) - - fun 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) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/SessionBackend.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/SessionBackend.kt index 1119d9e9be..f379364cc4 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/SessionBackend.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/SessionBackend.kt @@ -78,9 +78,9 @@ interface SessionBackend { ctx: RequestContext, filters: List, 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( diff --git a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/prodbench/SmallReqFloorBenchmark.kt b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/prodbench/SmallReqFloorBenchmark.kt index cd6eb4525b..4c2c839965 100644 --- a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/prodbench/SmallReqFloorBenchmark.kt +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/prodbench/SmallReqFloorBenchmark.kt @@ -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() + var fanStart = 0L + val fanJobs = + (0 until FANOUT_SUBS).map { + val ready = CompletableDeferred() + 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(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()