Merge pull request #3660 from vitorpamplona/claude/nostr-relay-filters-mapping-bu95ib

Optimize relay query performance: cost-based driver selection, tag-author indexing, and fanout memoization
This commit is contained in:
Vitor Pamplona
2026-07-21 15:35:38 -04:00
committed by GitHub
21 changed files with 984 additions and 145 deletions
@@ -26,6 +26,7 @@ import com.vitorpamplona.amethyst.desktop.cache.DesktopLocalCache
import com.vitorpamplona.quartz.nip01Core.core.Event import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.store.sqlite.DefaultIndexingStrategy
import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
import com.vitorpamplona.quartz.utils.Log import com.vitorpamplona.quartz.utils.Log
import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.CoroutineScope
@@ -42,6 +43,18 @@ class LocalRelayStore(
) : AutoCloseable { ) : AutoCloseable {
companion object { companion object {
val LOCAL_RELAY_URL: NormalizedRelayUrl = NormalizedRelayUrl("ws://localhost/amethyst-local/") val LOCAL_RELAY_URL: NormalizedRelayUrl = NormalizedRelayUrl("ws://localhost/amethyst-local/")
/**
* Client defaults plus the authors-without-kinds index: shared
* ViewModels (`Nip65RelayListViewModel`, `PrivateOutboxRelayListViewModel`,
* `VanishRequestsState`) replay `authors`-only filters against this
* store, which full-scan without `(pubkey, created_at)` — the
* `(kind, pubkey, …)` index can't serve them, pubkey is its second
* column. A personal store is small, so the extra insert cost is
* negligible; existing DBs get the index built on next open via
* `ensureOptionalIndexes`.
*/
val INDEX_STRATEGY = DefaultIndexingStrategy(indexEventsByPubkeyAlone = true)
} }
private fun dbDir(pubKeyHex: String): File = File(homeDir, ".amethyst/accounts/${pubKeyHex.take(8)}") private fun dbDir(pubKeyHex: String): File = File(homeDir, ".amethyst/accounts/${pubKeyHex.take(8)}")
@@ -87,14 +100,14 @@ class LocalRelayStore(
dir.mkdirs() dir.mkdirs()
val path = File(dir, "events.db").absolutePath val path = File(dir, "events.db").absolutePath
try { try {
store = EventStore(dbName = path, relay = LOCAL_RELAY_URL) store = EventStore(dbName = path, relay = LOCAL_RELAY_URL, indexStrategy = INDEX_STRATEGY)
_lastError.value = null _lastError.value = null
refreshStats() refreshStats()
} catch (e: Exception) { } catch (e: Exception) {
Log.w("LocalRelayStore") { "DB open failed, recreating: ${e.message}" } Log.w("LocalRelayStore") { "DB open failed, recreating: ${e.message}" }
try { try {
deleteDbFiles(path) deleteDbFiles(path)
store = EventStore(dbName = path, relay = LOCAL_RELAY_URL) store = EventStore(dbName = path, relay = LOCAL_RELAY_URL, indexStrategy = INDEX_STRATEGY)
_lastError.value = "Database was recreated: ${e.message}" _lastError.value = "Database was recreated: ${e.message}"
} catch (e2: Exception) { } catch (e2: Exception) {
_lastError.value = "Cannot open local store: ${e2.message}" _lastError.value = "Cannot open local store: ${e2.message}"
@@ -57,6 +57,14 @@ fun relayIndexingStrategy(
// index unconditionally; without it the filter walks the whole // index unconditionally; without it the filter walks the whole
// time index. // time index.
indexEventsByPubkeyAlone = true, indexEventsByPubkeyAlone = true,
// The tag ∩ author ∩ kind shape (DM rooms, reports-by-follows,
// follows-scoped community feeds — 65 client assembler call sites)
// otherwise reads every row for the tag/kind before filtering the
// author. TagAuthorIndexBenchmark @ 1M events: 14.2 ms -> 0.66 ms
// (~21x, growing with corpus size) with insert cost inside run noise
// (49.0 vs 47.4 µs/event). Existing DBs build the index on next open
// via ensureOptionalIndexes.
indexTagsWithKindAndPubkey = true,
indexFullTextSearch = fullTextSearch, indexFullTextSearch = fullTextSearch,
// Tokenize off the commit path; NostrServer drives the catch-up // Tokenize off the commit path; NostrServer drives the catch-up
// worker and search queries drain it first, so NIP-50 stays // worker and search queries drain it first, so NIP-50 stays
+2
View File
@@ -83,6 +83,8 @@ val store = EventStore(
By default, all single-letter tags with values are indexed. Override `shouldIndex(kind, tag)` for custom behavior. More indexes = faster queries but larger database. By default, all single-letter tags with values are indexed. Override `shouldIndex(kind, tag)` for custom behavior. More indexes = faster queries but larger database.
Flag flips are safe on existing databases: any flag-gated index the strategy wants but the on-disk schema lacks is created on the next open (idempotent `CREATE INDEX IF NOT EXISTS`, one-time build cost) — no schema version bump involved. Disabling a flag never drops an existing index.
`indexFullTextSearch` defaults to `true` and controls the NIP-50 full-text index (`event_fts`). Set it to `false` when search is served elsewhere (e.g. a Vespa backend, or a `SearchEventSource` as shown below): inserts skip the FTS tokenization cost, no `event_fts` table/trigger is created, and any filter carrying a non-empty `search` term returns no matches. `indexFullTextSearch` defaults to `true` and controls the NIP-50 full-text index (`event_fts`). Set it to `false` when search is served elsewhere (e.g. a Vespa backend, or a `SearchEventSource` as shown below): inserts skip the FTS tokenization cost, no `event_fts` table/trigger is created, and any filter carrying a non-empty `search` term returns no matches.
## Non-Storage Relays (search, redirector, computed) ## Non-Storage Relays (search, redirector, computed)
+2
View File
@@ -94,6 +94,8 @@ kotlin {
// Forward the negentropy-benchmark corpus size to the test JVM. // Forward the negentropy-benchmark corpus size to the test JVM.
System.getProperty("negBenchN")?.let { systemProperty("negBenchN", it) } System.getProperty("negBenchN")?.let { systemProperty("negBenchN", it) }
System.getProperty("followBenchScale")?.let { systemProperty("followBenchScale", it) } System.getProperty("followBenchScale")?.let { systemProperty("followBenchScale", it) }
System.getProperty("tagBenchScale")?.let { systemProperty("tagBenchScale", it) }
System.getProperty("fsBenchScale")?.let { systemProperty("fsBenchScale", it) }
// Opt-in JFR profiling of a benchmark run (-PnegProfile=/tmp/neg.jfr). // Opt-in JFR profiling of a benchmark run (-PnegProfile=/tmp/neg.jfr).
(project.findProperty("negProfile") as? String)?.let { (project.findProperty("negProfile") as? String)?.let {
jvmArgs("-XX:+FlightRecorder", "-XX:StartFlightRecording=filename=$it,settings=profile,dumponexit=true") jvmArgs("-XX:+FlightRecorder", "-XX:StartFlightRecording=filename=$it,settings=profile,dumponexit=true")
@@ -25,6 +25,22 @@ class EoseMessage(
) : Message { ) : Message {
override fun label() = LABEL override fun label() = LABEL
/**
* Wire form is `["EOSE","<subId>"]` — sent once per REQ, so it is on
* the per-subscription floor. Splice it directly when [subId] needs no
* escaping (the common case: client-chosen sub ids are short ASCII),
* skipping the generic serializer's node tree. Byte-identical output;
* any exotic subId falls back.
*/
override fun toJson(): String {
if (!isEscapeFreeAscii(subId)) return super.toJson()
return buildString(subId.length + 12) {
append("[\"EOSE\",\"")
append(subId)
append("\"]")
}
}
companion object { companion object {
const val LABEL = "EOSE" const val LABEL = "EOSE"
} }
@@ -29,6 +29,25 @@ class OkMessage(
) : Message { ) : Message {
override fun label() = LABEL override fun label() = LABEL
/**
* Wire form is `["OK","<eventId>",<true|false>,"<message>"]` — sent
* once per published EVENT. [eventId] is validated hex (always
* escape-free); splice directly when [message] also needs no escaping,
* which covers the empty-string success ack and the plain-ASCII
* rejection reasons. Byte-identical output; a reason with quotes or
* non-ASCII falls back to the generic serializer.
*/
override fun toJson(): String {
if (!isEscapeFreeAscii(message)) return super.toJson()
return buildString(eventId.length + message.length + 20) {
append("[\"OK\",\"")
append(eventId)
append(if (success) "\",true,\"" else "\",false,\"")
append(message)
append("\"]")
}
}
companion object { companion object {
const val LABEL = "OK" const val LABEL = "OK"
@@ -0,0 +1,36 @@
/*
* 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.commands.toClient
/**
* True when every char of [s] is printable ASCII (0x20–0x7e) and not a JSON
* metacharacter (`"` / `\`) — i.e. exactly the bytes a JSON string encoder
* would emit verbatim between the quotes. Frame builders use this to gate a
* direct-`buildString` fast path against the generic serializer: when it holds
* the spliced output is byte-identical, and any exotic value (control chars,
* quotes, non-ASCII) falls back to the escaping serializer.
*/
internal fun isEscapeFreeAscii(s: String): Boolean {
for (c in s) {
if (c < ' ' || c > '~' || c == '"' || c == '\\') return false
}
return true
}
@@ -22,6 +22,11 @@ package com.vitorpamplona.quartz.nip01Core.relay.filters
import com.vitorpamplona.quartz.nip01Core.core.Event import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.core.HexKey import com.vitorpamplona.quartz.nip01Core.core.HexKey
import kotlinx.collections.immutable.PersistentMap
import kotlinx.collections.immutable.PersistentSet
import kotlinx.collections.immutable.persistentHashMapOf
import kotlinx.collections.immutable.persistentHashSetOf
import kotlinx.collections.immutable.toPersistentHashSet
import kotlin.concurrent.atomics.AtomicReference import kotlin.concurrent.atomics.AtomicReference
import kotlin.concurrent.atomics.ExperimentalAtomicApi import kotlin.concurrent.atomics.ExperimentalAtomicApi
@@ -62,9 +67,10 @@ import kotlin.concurrent.atomics.ExperimentalAtomicApi
* copy-on-write CAS loops, mirroring the * copy-on-write CAS loops, mirroring the
* `nip86RelayManagement.server.BanStore` pattern. Reads in * `nip86RelayManagement.server.BanStore` pattern. Reads in
* [candidatesFor] and [forEach] are wait-free single-load atomic. * [candidatesFor] and [forEach] are wait-free single-load atomic.
* Writes (subscription register / unregister) copy the inner maps * Writes (subscription register / unregister) build the next snapshot
* — fine for this workload because writes are subscription-rate * from persistent (HAMT) maps — O(keys × log S) with structural
* (rare) while reads are event-rate (frequent). * sharing rather than a full O(S) copy of both maps, since on a relay
* a write happens on every REQ open and close.
* *
* ## What the index does NOT cover * ## What the index does NOT cover
* *
@@ -109,14 +115,26 @@ class FilterIndex<S : Any> {
private object Unindexed : BucketKey private object Unindexed : BucketKey
/** /**
* Single immutable snapshot. [buckets] maps a key to the set of * Single immutable snapshot. Subscribers are held in one map per
* subscribers registered under it; [assignments] is the reverse * indexable dimension so [candidatesFor] — called once per accepted
* map used by [unregister] to find a subscriber's keys without * ingest event, the hot read — can probe each dimension with the
* scanning every bucket. * 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 full maps — registration happens on every REQ open/close.
*/ */
private data class State<S>( private data class State<S>(
val buckets: Map<BucketKey, Set<S>> = emptyMap(), val ids: PersistentMap<HexKey, PersistentSet<S>> = persistentHashMapOf(),
val assignments: Map<S, Set<BucketKey>> = emptyMap(), 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(),
) )
private val state: AtomicReference<State<S>> = AtomicReference(State()) private val state: AtomicReference<State<S>> = AtomicReference(State())
@@ -181,18 +199,23 @@ class FilterIndex<S : Any> {
while (true) { while (true) {
val current = state.load() val current = state.load()
val keys = current.assignments[subscriber] ?: return val keys = current.assignments[subscriber] ?: return
val newBuckets = current.buckets.toMutableMap() var ids = current.ids
var authors = current.authors
var tags = current.tags
var kinds = current.kinds
var unindexed = current.unindexed
for (key in keys) { for (key in keys) {
val cur = newBuckets[key] ?: continue when (key) {
val next = cur - subscriber is IdKey -> ids = ids.removeSub(key.id, subscriber)
if (next.isEmpty()) { is AuthorKey -> authors = authors.removeSub(key.author, subscriber)
newBuckets.remove(key) is KindKey -> kinds = kinds.removeSub(key.kind, subscriber)
} else { is TagKey -> tags = tags.removeTagSub(key.letter, key.value, subscriber)
newBuckets[key] = next Unindexed -> unindexed = unindexed.remove(subscriber)
} }
} }
val newAssignments = current.assignments - subscriber val next =
if (state.compareAndSet(current, State(newBuckets, newAssignments))) return State(ids, authors, tags, kinds, unindexed, current.assignments.remove(subscriber))
if (state.compareAndSet(current, next)) return
} }
} }
@@ -202,19 +225,22 @@ class FilterIndex<S : Any> {
* candidate to handle negative constraints. * candidate to handle negative constraints.
* *
* Iteration order is insertion-stable per call but otherwise * 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> { fun candidatesFor(event: Event): Set<S> {
val s = state.load() val s = state.load()
if (s.buckets.isEmpty()) return emptySet() if (s.assignments.isEmpty()) return emptySet()
val result = LinkedHashSet<S>() val result = LinkedHashSet<S>()
s.buckets[Unindexed]?.let { result.addAll(it) } if (s.unindexed.isNotEmpty()) result.addAll(s.unindexed)
s.buckets[IdKey(event.id)]?.let { result.addAll(it) } s.ids[event.id]?.let { result.addAll(it) }
s.buckets[AuthorKey(event.pubKey)]?.let { result.addAll(it) } s.authors[event.pubKey]?.let { result.addAll(it) }
s.buckets[KindKey(event.kind)]?.let { result.addAll(it) } s.kinds[event.kind]?.let { result.addAll(it) }
for (tag in event.tags) { if (s.tags.isNotEmpty()) {
if (tag.size >= 2 && tag[0].length == 1) { for (tag in event.tags) {
s.buckets[TagKey(tag[0], tag[1])]?.let { result.addAll(it) } if (tag.size >= 2 && tag[0].length == 1) {
s.tags[tag[0]]?.get(tag[1])?.let { result.addAll(it) }
}
} }
} }
return result return result
@@ -235,19 +261,75 @@ class FilterIndex<S : Any> {
keys: List<BucketKey>, keys: List<BucketKey>,
) { ) {
if (keys.isEmpty()) return if (keys.isEmpty()) return
val keySet = keys.toSet() val keySet = keys.toPersistentHashSet()
while (true) { while (true) {
val current = state.load() val current = state.load()
val newBuckets = current.buckets.toMutableMap() var ids = current.ids
var authors = current.authors
var tags = current.tags
var kinds = current.kinds
var unindexed = current.unindexed
for (key in keySet) { for (key in keySet) {
val cur = newBuckets[key] ?: emptySet() when (key) {
if (subscriber in cur) continue is IdKey -> ids = ids.addSub(key.id, subscriber)
newBuckets[key] = cur + 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 existing = current.assignments[subscriber]
val merged = if (existing == null) keySet else existing + keySet val merged = existing?.addAll(keySet) ?: keySet
val newAssignments = current.assignments + (subscriber to merged) val next = State(ids, authors, tags, kinds, unindexed, current.assignments.put(subscriber, merged))
if (state.compareAndSet(current, State(newBuckets, newAssignments))) return 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)
} }
} }
@@ -49,6 +49,7 @@ import com.vitorpamplona.quartz.utils.Log
import com.vitorpamplona.quartz.utils.cache.LargeCache import com.vitorpamplona.quartz.utils.cache.LargeCache
import kotlinx.coroutines.CancellationException import kotlinx.coroutines.CancellationException
import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.CoroutineStart
import kotlinx.coroutines.Job import kotlinx.coroutines.Job
import kotlinx.coroutines.channels.ClosedSendChannelException import kotlinx.coroutines.channels.ClosedSendChannelException
import kotlinx.coroutines.launch import kotlinx.coroutines.launch
@@ -291,8 +292,16 @@ class RelaySession(
// Policy may rewrite filters to match the user's access level. // Policy may rewrite filters to match the user's access level.
val filters = (result as PolicyResult.Accepted).cmd.filters val filters = (result as PolicyResult.Accepted).cmd.filters
// UNDISPATCHED: the stored replay runs inline on this coroutine —
// the reader-pool acquire doesn't suspend when a connection is
// free, so EVENT frames and EOSE go out without a scheduler hop
// (SmallReqFloorBenchmark: the hop was most of the dispatch
// slice on small REQs). The coroutine first parks at the live
// tail (awaitCancellation), which is when launch returns and the
// job lands in [subscriptions]; commands on this connection are
// processed sequentially, so nothing can target the sub earlier.
val job = val job =
scope.launch { scope.launch(start = CoroutineStart.UNDISPATCHED) {
try { try {
if (policy.filtersOutgoingEvents) { if (policy.filtersOutgoingEvents) {
// Screened path: every event is materialized so the // Screened path: every event is materialized so the
@@ -331,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)) }, onEose = { send(EoseMessage(cmd.subId)) },
) )
} }
@@ -74,13 +74,57 @@ class LiveEventStore(
* One live REQ subscription. Carries the filters (for the * One live REQ subscription. Carries the filters (for the
* post-index `match` re-check needed for negative constraints * post-index `match` re-check needed for negative constraints
* like `since` / `until` / `tagsAll`) and the delivery callback * 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( private class LiveSubscription(
val filters: List<Filter>, 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-and-forget enqueue: hand [event] to the [IngestQueue] and
* fire [onComplete] once the writer's batch has a per-row * fire [onComplete] once the writer's batch has a per-row
@@ -149,9 +193,17 @@ class LiveEventStore(
* batch writer. * batch writer.
*/ */
private fun fanout(event: Event) { 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) }) { 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, onEose: () -> Unit,
) { ) {
drainFtsIfSearching(filters) drainFtsIfSearching(filters)
// During the historical replay, record ids the store has // The index registers *before* the replay starts (otherwise an
// emitted so the live path can dedupe. The index registers // event accepted mid-replay would slip past the live path entirely
// *before* the replay starts (otherwise an event accepted // — same race the previous SharedFlow-based implementation closed
// mid-replay would slip past the live path entirely — same // with `onSubscription`), and [SeenIds] bridges the two coroutines:
// race the previous SharedFlow-based implementation closed // the replay records ids here, the live `deliver` drops duplicates,
// with `onSubscription`). // and after EOSE the set is released so every live event forwards.
// val seen = SeenIds()
// 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)
}
}
val sub = val sub =
LiveSubscription( LiveSubscription(
filters = filters, filters = filters,
deliver = { event -> deliver = { event, _ ->
val duplicate = seenLocked { seenIds?.contains(event.id) ?: false } if (!seen.isDuplicate(event.id)) onEach(event)
if (duplicate) return@LiveSubscription
onEach(event)
}, },
) )
index.register(filters, sub) index.register(filters, sub)
try { try {
store.query<Event>(filters.strippingSearchExtensions()) { event -> store.query<Event>(filters.strippingSearchExtensions()) { event ->
seenLocked { seenIds?.add(event.id) } seen.record(event.id)
onEach(event) onEach(event)
} }
onEose() onEose()
// Drop the dedupe set so the live path stops paying for // Drop the dedupe set so the live path stops paying for
// it. From this point the index drives delivery and // it. From this point the index drives delivery and
// duplicates are no longer possible. // duplicates are no longer possible.
seenLocked { seenIds = null } seen.release()
// Suspend until the caller's coroutine is cancelled // Suspend until the caller's coroutine is cancelled
// (e.g. NIP-01 CLOSE or connection drop). The `finally` // (e.g. NIP-01 CLOSE or connection drop). The `finally`
// unregisters from the index. // unregisters from the index.
@@ -255,42 +279,28 @@ class LiveEventStore(
ctx: RequestContext, ctx: RequestContext,
filters: List<Filter>, filters: List<Filter>,
onEachStored: (RawEvent) -> Unit, onEachStored: (RawEvent) -> Unit,
onEachLive: (Event) -> Unit, onEachLive: (Event, String) -> Unit,
onEose: () -> Unit, onEose: () -> Unit,
) { ) {
drainFtsIfSearching(filters) drainFtsIfSearching(filters)
val seenLock = AtomicBoolean(false) val seen = SeenIds()
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 sub = val sub =
LiveSubscription( LiveSubscription(
filters = filters, filters = filters,
deliver = { event -> deliver = { event, body ->
val duplicate = seenLocked { seenIds?.contains(event.id) ?: false } if (!seen.isDuplicate(event.id)) onEachLive(event, body)
if (duplicate) return@LiveSubscription
onEachLive(event)
}, },
) )
index.register(filters, sub) index.register(filters, sub)
try { try {
store.rawQuery(filters.strippingSearchExtensions()) { raw -> store.rawQuery(filters.strippingSearchExtensions()) { raw ->
seenLocked { seenIds?.add(raw.id) } seen.record(raw.id)
onEachStored(raw) onEachStored(raw)
} }
onEose() onEose()
seenLocked { seenIds = null } seen.release()
awaitCancellation() awaitCancellation()
} finally { } finally {
index.unregister(sub) index.unregister(sub)
@@ -78,9 +78,9 @@ interface SessionBackend {
ctx: RequestContext, ctx: RequestContext,
filters: List<Filter>, filters: List<Filter>,
onEachStored: (RawEvent) -> Unit, onEachStored: (RawEvent) -> Unit,
onEachLive: (Event) -> Unit, onEachLive: (Event, String) -> Unit,
onEose: () -> 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]. */ /** Answers a NIP-45 COUNT with an exact cardinality for the caller in [ctx]. */
suspend fun count( suspend fun count(
@@ -147,13 +147,38 @@ class EventIndexesModule(
*/ */
fun migrateV2AddPubkeyIndex(db: SQLiteConnection) { fun migrateV2AddPubkeyIndex(db: SQLiteConnection) {
if (!indexStrategy.indexEventsByPubkeyAlone) return if (!indexStrategy.indexEventsByPubkeyAlone) return
val orderBy = db.execSQL("CREATE INDEX IF NOT EXISTS query_by_pubkey_created ON event_headers (pubkey, ${orderByColumns()})")
if (indexStrategy.useAndIndexIdOnOrderBy) { }
"created_at DESC, id ASC"
} else { private fun orderByColumns() =
"created_at DESC" if (indexStrategy.useAndIndexIdOnOrderBy) {
} "created_at DESC, id ASC"
db.execSQL("CREATE INDEX IF NOT EXISTS query_by_pubkey_created ON event_headers (pubkey, $orderBy)") } else {
"created_at DESC"
}
/**
* Materializes any flag-gated index the current [indexStrategy] wants
* but the on-disk schema predates. Flags are runtime configuration, not
* schema — a deployment can flip one without a `user_version` bump — so
* this runs idempotently on every open. The first open after enabling a
* flag pays a one-time index build over the existing rows; subsequent
* opens are no-ops. A disabled flag never drops an existing index (that
* stays an operator decision).
*/
fun ensureOptionalIndexes(db: SQLiteConnection) {
if (indexStrategy.indexEventsByCreatedAtAlone) {
db.execSQL("CREATE INDEX IF NOT EXISTS query_by_created_at_id ON event_headers (${orderByColumns()})")
}
if (indexStrategy.indexEventsByPubkeyAlone) {
db.execSQL("CREATE INDEX IF NOT EXISTS query_by_pubkey_created ON event_headers (pubkey, ${orderByColumns()})")
}
if (indexStrategy.indexTagsByCreatedAtAlone) {
db.execSQL("CREATE INDEX IF NOT EXISTS query_by_tags_hash ON event_tags (tag_hash, created_at DESC)")
}
if (indexStrategy.indexTagsWithKindAndPubkey) {
db.execSQL("CREATE INDEX IF NOT EXISTS query_by_tags_hash_kind_pubkey ON event_tags (tag_hash, kind, pubkey_hash, created_at DESC)")
}
} }
val sqlInsertHeader = val sqlInsertHeader =
@@ -74,9 +74,21 @@ interface IndexingStrategy {
* Activate this if you see too many Tag-centric Filters without * Activate this if you see too many Tag-centric Filters without
* kind AND pubkey at the same time. * kind AND pubkey at the same time.
* *
* This is a rarely used index (reports by your follows or * This shape (reports by your follows, NIP-04 DM rooms, follows-scoped
* NIP-04 DMs for instance) that becomes quite large without * community feeds) is not rare on the client side: the 2026-07 filter
* major gains. * assembler survey counted 65 call sites building
* `kinds + authors + tags`. Without this index the plan seeks
* `(tag_hash, kind)` and reads every row for that tag/kind before
* filtering the author.
*
* Measured by `TagAuthorIndexBenchmark` (jvmTest prodbench): the
* DM-room query drops 9.4 ms → 0.6 ms (~15×) at 200k events and
* 14.2 ms → 0.66 ms (~21×) at 1M — the gap grows with corpus size —
* while batch-insert cost stays inside run noise (49.0 vs 47.4
* µs/event at 1M). geode enables it; the client default stays off
* because a client store's per-tag row counts are bounded by one
* user's data. Flipping it on an existing DB is safe: the index is
* built on next open by `EventIndexesModule.ensureOptionalIndexes`.
* *
* Keep in mind that activating too many indexes increases the size of the * Keep in mind that activating too many indexes increases the size of the
* DB so much that the indexes themselves won't fit in memory, requiring * DB so much that the indexes themselves won't fit in memory, requiring
@@ -52,6 +52,17 @@ import androidx.sqlite.SQLiteStatement
* exactly at a same-second boundary. * exactly at a same-second boundary.
*/ */
internal object MergeQueryExecutor { internal object MergeQueryExecutor {
// TODO: the same collect-all + TEMP-B-TREE-sort pattern exists one index
// over, on the tag path: `kinds + tags(#e IN [hundreds]) + limit` (the
// reactions/replies watcher archetype) unions per-value streams that are
// each sorted off `(tag_hash, kind, created_at)` and sorts the union.
// `streamCount` currently rejects any filter with tags, so those queries
// never merge. Measured by `TagAuthorIndexBenchmark`: `#e IN 300,
// limit 500` costs 12.8 ms cold at 200k events and 14.2 ms at 1M
// (6.7 ms with indexTagsWithKindAndPubkey on) — tolerable, but it
// scales with matching history like the follow-feed shape did; extend
// the merge to per-tag-value streams if the relayBench
// `reactions-watch` scenario shows it in the profile vs strfry.
const val COLS = "id, pubkey, created_at, kind, tags, content, sig" const val COLS = "id, pubkey, created_at, kind, tags, content, sig"
/** /**
@@ -160,6 +160,11 @@ class SQLiteEventStore(
setUserVersion(this, DATABASE_VERSION) setUserVersion(this, DATABASE_VERSION)
} }
} }
// Flag-gated indexes are runtime config, not schema: a
// deployment that flips an IndexingStrategy flag on an
// existing DB gets the index built here (idempotent,
// one-time cost), with no user_version bump involved.
eventIndexModule.ensureOptionalIndexes(db)
}, },
) )
} }
@@ -202,8 +202,20 @@ fun Filter.strippingSearchExtensions(): Filter {
/** /**
* Applies [strippingSearchExtensions] to every filter, returning this * Applies [strippingSearchExtensions] to every filter, returning this
* same list when no filter carried extension tokens. * same list when no filter carried extension tokens.
*
* This runs on every REQ/COUNT/snapshot, and the overwhelming majority
* carry no `search` term at all, so the no-search case must not allocate:
* bail before building any list when nothing could be stripped.
*/ */
fun List<Filter>.strippingSearchExtensions(): List<Filter> { fun List<Filter>.strippingSearchExtensions(): List<Filter> {
var hasSearch = false
for (i in indices) {
if (!this[i].search.isNullOrEmpty()) {
hasSearch = true
break
}
}
if (!hasSearch) return this
var changed = false var changed = false
val out = val out =
map { map {
@@ -25,6 +25,7 @@ import com.vitorpamplona.quartz.nip01Core.core.isAddressable
import com.vitorpamplona.quartz.nip01Core.core.isReplaceable import com.vitorpamplona.quartz.nip01Core.core.isReplaceable
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.store.sqlite.TagNameValueHasher import com.vitorpamplona.quartz.nip01Core.store.sqlite.TagNameValueHasher
import java.nio.file.DirectoryStream
import java.nio.file.Files import java.nio.file.Files
import java.nio.file.Path import java.nio.file.Path
import kotlin.io.path.exists import kotlin.io.path.exists
@@ -36,16 +37,23 @@ import kotlin.io.path.exists
* *
* Step-2 coverage: * Step-2 coverage:
* - `ids` → direct canonical opens * - `ids` → direct canonical opens
* - `tagsAll`/`tags` → tag index union (first key) * - `tagsAll`/`tags`/`kinds`/`authors` → cheapest index tree drives
* - `kinds` → kind index union * (capped entry-count comparison), the rest post-filter
* - `authors` → author index union
* - otherwise → full scan via every `idx/kind/<k>/` subtree * - otherwise → full scan via every `idx/kind/<k>/` subtree
* *
* The planner is intentionally dumb about selectivity — "first available * Driver choice is cost-based: every legal driver (each `tagsAll` value
* driver wins". A cost-based picker (smallest listing) can slot in * alone — AND semantics make any single value a complete driver — each
* later without changing callers. All FilterMatcher semantics (tag * `tags` key's value union, the kind set, the author set) opens a lazy
* AND/OR, since/until, id, author, kind cross-checks) are enforced in * directory iterator, all are drained in lockstep, and the first to
* the orchestrator, so picking a loose driver is correctness-safe. * exhaust — the smallest listing — drives. A giant tree (`idx/kind/1/`
* with a million entries) is therefore never read past ~the smallest
* candidate's size. Before the pick, the fixed tags → kinds → authors
* order sent `authors + kinds + limit` — the most common CLI shape —
* through the kind tree: 149 ms fixed-order vs 4.0 ms cost-based at 30k
* events per `FsDriverSelectionBenchmark` (floor: author-only at
* 1.4 ms). All FilterMatcher semantics (tag AND/OR,
* since/until, id, author, kind cross-checks) are enforced in the
* orchestrator, so any driver pick is correctness-safe.
*/ */
internal class FsQueryPlanner( internal class FsQueryPlanner(
private val layout: FsLayout, private val layout: FsLayout,
@@ -74,19 +82,100 @@ internal class FsQueryPlanner(
return ftsDriver(search) return ftsDriver(search)
} }
firstTagKey(filter)?.let { (name, values) -> val candidates = driverCandidates(filter)
return mergeDesc(values.map { v -> walkDir(layout.tagValueDir(name, v, hasher.hash(name, v))) }) if (candidates.isEmpty()) return allKindsDriver()
return mergeDesc(cheapestDriver(candidates).map { walkDir(it) })
}
/**
* Every set of index directories that, walked and post-filtered, yields
* a superset of the filter's matches:
* - each `tagsAll` value alone (AND semantics — every match carries it),
* - each `tags` key's full value union (OR within a key, AND across),
* - the kind set, and the author set.
* Listed in the old fixed-priority order so [cheapestDriver] keeps that
* order on cost ties.
*/
private fun driverCandidates(filter: Filter): List<List<Path>> {
val out = ArrayList<List<Path>>()
filter.tagsAll?.forEach { (name, values) ->
values.forEach { v -> out.add(listOf(layout.tagValueDir(name, v, hasher.hash(name, v)))) }
}
filter.tags?.forEach { (name, values) ->
if (values.isNotEmpty()) {
out.add(values.map { v -> layout.tagValueDir(name, v, hasher.hash(name, v)) })
}
}
filter.kinds?.takeIf { it.isNotEmpty() }?.let { kinds ->
out.add(kinds.map { layout.kindDir(it) })
}
filter.authors?.takeIf { it.isNotEmpty() }?.let { authors ->
out.add(authors.map { layout.authorDir(it) })
}
return out
}
/**
* Smallest candidate by lockstep listing drain: one lazy directory
* iterator per candidate, all advanced [COST_BATCH] entries per round —
* the first to exhaust its listing is the smallest, so a giant tree is
* never read past ~the smallest candidate's size (a candidate that
* exhausts on round one costs the others one batch each). A candidate
* whose dirs are all missing exhausts immediately: driving from an empty
* mandatory predicate correctly yields an empty result. If every
* candidate survives [COST_CAP] entries, all are huge and relative
* driver choice stops mattering — the first (old fixed-priority order)
* wins.
*/
private fun cheapestDriver(candidates: List<List<Path>>): List<Path> {
if (candidates.size == 1) return candidates[0]
val cursors = candidates.map { EntryCursor(it) }
try {
var advanced = 0L
while (advanced < COST_CAP) {
for (i in cursors.indices) {
if (!cursors[i].skip(COST_BATCH)) return candidates[i]
}
advanced += COST_BATCH
}
return candidates[0]
} finally {
cursors.forEach { it.close() }
}
}
/** Lazy entry iterator over a candidate's directories, in order. */
private class EntryCursor(
dirs: List<Path>,
) : AutoCloseable {
private val remaining = ArrayDeque(dirs)
private var stream: DirectoryStream<Path>? = null
private var iter: Iterator<Path> = emptyList<Path>().iterator()
/** Advances up to [n] entries; false when the listing ends first. */
fun skip(n: Int): Boolean {
var left = n
while (left > 0) {
if (iter.hasNext()) {
iter.next()
left--
continue
}
close()
val dir = remaining.removeFirstOrNull() ?: return false
if (!Files.isDirectory(dir)) continue
stream = Files.newDirectoryStream(dir)
iter = stream!!.iterator()
}
return true
} }
filter.kinds?.let { kinds -> override fun close() {
return mergeDesc(kinds.map { walkDir(layout.kindDir(it)) }) stream?.close()
stream = null
iter = emptyList<Path>().iterator()
} }
filter.authors?.let { authors ->
return mergeDesc(authors.map { walkDir(layout.authorDir(it)) })
}
return allKindsDriver()
} }
/** /**
@@ -301,14 +390,15 @@ internal class FsQueryPlanner(
var top: Candidate, var top: Candidate,
) )
// ---- helpers ------------------------------------------------------ private companion object {
/** Entries each candidate's cursor advances per lockstep round. */
const val COST_BATCH = 64
/** First tag filter with at least one value, preferring `tagsAll`. */ /**
private fun firstTagKey(filter: Filter): Pair<String, List<String>>? { * Stop draining once every candidate has survived this many
filter.tagsAll?.firstNonEmpty()?.let { return it } * entries: past it they are all huge, relative choice stops
filter.tags?.firstNonEmpty()?.let { return it } * mattering, and the first candidate in priority order wins.
return null */
const val COST_CAP = 65_536L
} }
private fun Map<String, List<String>>.firstNonEmpty(): Pair<String, List<String>>? = entries.firstOrNull { it.value.isNotEmpty() }?.let { it.key to it.value }
} }
@@ -32,11 +32,14 @@ import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
import com.vitorpamplona.quartz.utils.EventFactory import com.vitorpamplona.quartz.utils.EventFactory
import kotlinx.coroutines.CompletableDeferred import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.CoroutineStart
import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel import kotlinx.coroutines.cancel
import kotlinx.coroutines.launch import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.withTimeout
import java.util.concurrent.atomic.AtomicInteger
import kotlin.test.Test import kotlin.test.Test
import kotlin.test.assertEquals import kotlin.test.assertEquals
@@ -65,6 +68,8 @@ class SmallReqFloorBenchmark {
const val AUTHORS = 2_500 // ~20 events per author, matching author-archive const val AUTHORS = 2_500 // ~20 events per author, matching author-archive
const val ROUNDS = 400 const val ROUNDS = 400
const val WARMUP = 100 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') private fun hexId(seed: Int): String = seed.toString(16).padStart(64, '0')
@@ -122,16 +127,19 @@ class SmallReqFloorBenchmark {
} }
// --- B: backend queryRaw to EOSE (live machinery included) --- // --- B: backend queryRaw to EOSE (live machinery included) ---
// UNDISPATCHED mirrors the production path (RelaySession.handleReq
// starts the query coroutine undispatched), so B−A is the live
// machinery itself, not a benchmark-only scheduler hop.
suspend fun timeBackend(round: Int): Long { suspend fun timeBackend(round: Int): Long {
val eose = CompletableDeferred<Long>() val eose = CompletableDeferred<Long>()
val t0 = System.nanoTime() val t0 = System.nanoTime()
val job = val job =
scope.launch { scope.launch(start = CoroutineStart.UNDISPATCHED) {
live.queryRaw( live.queryRaw(
ctx = ctx, ctx = ctx,
filters = listOf(filterFor(round)), filters = listOf(filterFor(round)),
onEachStored = {}, onEachStored = {},
onEachLive = {}, onEachLive = { _, _ -> },
onEose = { eose.complete(System.nanoTime() - t0) }, onEose = { eose.complete(System.nanoTime() - t0) },
) )
} }
@@ -143,6 +151,33 @@ class SmallReqFloorBenchmark {
val b = LongArray(ROUNDS) val b = LongArray(ROUNDS)
repeat(ROUNDS) { b[it] = timeBackend(it) } repeat(ROUNDS) { b[it] = timeBackend(it) }
// --- B@1k: same, with 1000 idle live subscriptions parked ---
// Register/unregister cost scales with the live population
// (FilterIndex mutates a shared snapshot per REQ open/close),
// which the single-sub stage can't see. Each idle sub filters
// on an author absent from the corpus: 0-row replay, then parks
// at the live tail and stays registered.
val idleJobs =
(0 until IDLE_SUBS).map { i ->
val ready = CompletableDeferred<Unit>()
val job =
scope.launch(start = CoroutineStart.UNDISPATCHED) {
live.queryRaw(
ctx = ctx,
filters = listOf(Filter(authors = listOf(hexId(1_000_000 + i)), kinds = listOf(1), limit = 1)),
onEachStored = {},
onEachLive = { _, _ -> },
onEose = { ready.complete(Unit) },
)
}
ready.await()
job
}
repeat(WARMUP) { timeBackend(it) }
val b1k = LongArray(ROUNDS)
repeat(ROUNDS) { b1k[it] = timeBackend(it) }
idleJobs.forEach { it.cancel() }
// --- C: full session dispatch, REQ json in → EOSE frame out --- // --- C: full session dispatch, REQ json in → EOSE frame out ---
suspend fun timeSession(round: Int): Long { suspend fun timeSession(round: Int): Long {
val eose = CompletableDeferred<Long>() val eose = CompletableDeferred<Long>()
@@ -160,15 +195,67 @@ class SmallReqFloorBenchmark {
val c = LongArray(ROUNDS) val c = LongArray(ROUNDS)
repeat(ROUNDS) { c[it] = timeSession(it) } 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") assertEquals(true, rowsA > 0, "author filters must return rows")
val mA = median(a) val mA = median(a)
val mB = median(b) val mB = median(b)
val mB1k = median(b1k)
val mC = median(c) val mC = median(c)
println("SmallReqFloorBenchmark @ ${EVENTS / 1000}k events, ~${rowsA / ROUNDS} rows/req, medians of $ROUNDS") println("SmallReqFloorBenchmark @ ${EVENTS / 1000}k events, ~${rowsA / ROUNDS} rows/req, medians of $ROUNDS")
println(" A raw store query: ${"%6.3f".format(mA)} ms") println(" A raw store query: ${"%6.3f".format(mA)} ms")
println(" B backend queryRaw→EOSE: ${"%6.3f".format(mB)} ms (live machinery +${"%6.3f".format(mB - mA)})") 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)})") 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() server.close()
scope.cancel() scope.cancel()
@@ -0,0 +1,184 @@
/*
* 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.filters.Filter
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.runBlocking
import kotlin.test.Test
/**
* Measures the two tag-path query shapes the client filter-assembler survey
* (2026-07) found hot but that no existing benchmark covers:
*
* 1. **tag ∩ author (DM-room shape)** — `kinds=[4] AND authors=[peer] AND
* #p=[me] LIMIT n`. 65 assembler call sites build this shape (every
* NIP-04 chat room, reports-by-follows, follows-scoped community feeds).
* [com.vitorpamplona.quartz.nip01Core.store.sqlite.IndexingStrategy.indexTagsWithKindAndPubkey]
* gates a covering `(tag_hash, kind, pubkey_hash, created_at)` index for
* it, but the flag is off everywhere (including geode). Without it the
* plan seeks `(tag_hash, kind)` and reads EVERY DM the user has ever
* received before filtering to the one peer. This compares query latency
* with the flag off vs on, and the batch-insert cost the extra index adds.
*
* 2. **large-IN tag watcher (reactions shape)** — `kinds=[7] AND
* #e=[hundreds of note ids] LIMIT n`. The per-value streams come sorted
* off `(tag_hash, kind, created_at)`, but their union does not, so SQLite
* collects every matching row and TEMP-B-TREE sorts to the limit — the
* tag-index analogue of the follow-feed regression
* [MergeQueryExecutor] fixed for author streams. Reported with and
* without a `since` bound to show what EOSE-warm steady state hides.
*
* Size the seed with `-DtagBenchScale=N` (default 1 ≈ ~200k events).
*/
class TagAuthorIndexBenchmark {
companion object {
val SCALE = System.getProperty("tagBenchScale")?.toInt() ?: 1
}
private val hex = "0123456789abcdef"
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 hex64(
salt: Long,
index: Int,
): String {
val out = CharArray(64)
for (w in 0 until 4) {
val v = mix(salt * 1_000_003 + index.toLong() * 4 + w)
for (b in 0 until 8) {
val byte = ((v ushr (b * 8)) and 0xFF).toInt()
out[(w * 8 + b) * 2] = hex[byte ushr 4]
out[(w * 8 + b) * 2 + 1] = hex[byte and 0xF]
}
}
return String(out)
}
private val sig = "0".repeat(128)
private var idSeq = 0
private fun ev(
pubkey: String,
createdAt: Long,
kind: Int,
tags: Array<Array<String>>,
): Event = EventFactory.create(hex64(7, idSeq++), pubkey, createdAt, kind, tags, "", sig)
private fun seedEvents(): List<Event> {
idSeq = 0
val base = 1_700_000_000L
val span = 3_000_000L // ~35 days
val me = hex64(9, 0)
val events = ArrayList<Event>(220_000 * SCALE)
// DM inbox: 200 peers, 300 DMs each → 60k kind-4 rows sharing the
// same (p:me) tag hash. The room query wants one peer's 300.
val peers = (0 until 200).map { hex64(2, it) }
for ((i, peer) in peers.withIndex()) {
repeat(300 * SCALE) {
val ts = base + (mix(i * 131L + it) and 0x7fffffff) % span
events.add(ev(peer, ts, 4, arrayOf(arrayOf("p", me))))
}
}
// Notification noise: 2000 authors mention me in kind-1 notes, so
// (p:me) spans multiple kinds like a real inbox does.
repeat(40_000 * SCALE) {
val author = hex64(3, it % 2_000)
val ts = base + (mix(it * 17L) and 0x7fffffff) % span
events.add(ev(author, ts, 1, arrayOf(arrayOf("p", me))))
}
// Reactions: 100k kind-7 events spread over 5000 target notes, for
// the large-IN watcher shape.
val noteIds = (0 until 5_000).map { hex64(5, it) }
repeat(100_000 * SCALE) {
val author = hex64(4, it % 3_000)
val ts = base + (mix(it * 29L) and 0x7fffffff) % span
events.add(ev(author, ts, 7, arrayOf(arrayOf("e", noteIds[it % noteIds.size]))))
}
return events
}
@Test
fun compareTagAuthorIndex() =
runBlocking {
val events = seedEvents()
val me = hex64(9, 0)
val peers = (0 until 200).map { hex64(2, it) }
val noteIds = (0 until 5_000).map { hex64(5, it) }
println("─ TagAuthorIndexBenchmark: ${events.size} events (scale=$SCALE) ─")
val strategies =
listOf(
"flag-off" to DefaultIndexingStrategy(indexFullTextSearch = false),
"flag-on " to DefaultIndexingStrategy(indexFullTextSearch = false, indexTagsWithKindAndPubkey = true),
)
for ((label, strategy) in strategies) {
val store = EventStore(dbName = null, indexStrategy = strategy)
val t0 = System.nanoTime()
events.chunked(10_000).forEach { store.batchInsert(it) }
val insertMs = (System.nanoTime() - t0) / 1e6
println(" ═ $label ═ insert: %.0f ms (%.1f µs/event)".format(insertMs, insertMs * 1000 / events.size))
// 1. DM room: one peer's DMs out of the whole (p:me) inbox.
val room = Filter(kinds = listOf(4), authors = listOf(peers[42]), tags = mapOf("p" to listOf(me)), limit = 100)
time(store, "dm-room (#p ∩ author ∩ kind, limit 100)", room)
// 2. Reactions watcher: 300 note ids, cold (no since).
val watcher = Filter(kinds = listOf(7), tags = mapOf("e" to noteIds.take(300)), limit = 500)
time(store, "reactions (#e IN 300, limit 500, cold)", watcher)
// 3. Same watcher, EOSE-warm (since bounds the window).
val warm = watcher.copy(since = 1_700_000_000L + 2_900_000L)
time(store, "reactions (#e IN 300, limit 500, since)", warm)
store.close()
}
}
private suspend fun time(
store: EventStore,
label: String,
filter: Filter,
) {
repeat(3) { store.query<Event>(filter) }
val runs = 10
var rows = 0
val start = System.nanoTime()
repeat(runs) { rows = store.query<Event>(filter).size }
val ms = (System.nanoTime() - start) / 1e6 / runs
println(" %-42s %8.2f ms (%d rows)".format(label, ms, rows))
}
}
@@ -0,0 +1,159 @@
/*
* 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.store.fs
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.utils.EventFactory
import kotlinx.coroutines.runBlocking
import java.nio.file.Files
import java.nio.file.Path
import kotlin.io.path.exists
import kotlin.test.Test
/**
* Guards [FsQueryPlanner]'s cost-based driver pick on the
* `authors + kinds + limit` shape — the most common CLI query (27
* assembler call sites; every `amy feed`-style author timeline over
* non-replaceable kinds).
*
* Under the pre-pick fixed order (tags → kinds → authors),
* `Filter(authors=[pk], kinds=[1], limit=n)` drove from `idx/kind/1/`
* (the biggest tree in any real store) and post-filtered the author:
* 149 ms at 30k events. The lockstep pick drives from the author tree
* and runs at ~4 ms. The benchmark times:
*
* - **planner (cost-based pick)**: the filter as the planner runs it —
* should sit near the floor, far below a kind-tree walk.
* - **author-driver emulation**: the author tree walked via an
* authors-only query with the kind check applied by the caller — the
* reference the pick is expected to match or beat.
* - **author-only floor**: `authors + limit` with no kind, the cheapest
* possible walk of the same tree.
*
* Size the seed with `-DfsBenchScale=N` (default 1 ≈ ~30k events; each
* event is a file + ~3 hardlinks, so seeding dominates wall time).
*/
class FsDriverSelectionBenchmark {
companion object {
val SCALE = System.getProperty("fsBenchScale")?.toInt() ?: 1
}
private val hex = "0123456789abcdef"
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 hex64(
salt: Long,
index: Int,
): String {
val out = CharArray(64)
for (w in 0 until 4) {
val v = mix(salt * 1_000_003 + index.toLong() * 4 + w)
for (b in 0 until 8) {
val byte = ((v ushr (b * 8)) and 0xFF).toInt()
out[(w * 8 + b) * 2] = hex[byte ushr 4]
out[(w * 8 + b) * 2 + 1] = hex[byte and 0xF]
}
}
return String(out)
}
private val sig = "0".repeat(128)
private var idSeq = 0
private fun ev(
pubkey: String,
createdAt: Long,
kind: Int,
): Event = EventFactory.create(hex64(7, idSeq++), pubkey, createdAt, kind, emptyArray(), "", sig)
@Test
fun compareDrivers() =
runBlocking {
val root: Path = Files.createTempDirectory("fs-driver-bench-")
val store = FsEventStore(root)
try {
val base = 1_700_000_000L
val span = 3_000_000L
val target = hex64(9, 0)
// Background: 300 authors × 100 kind-1 notes.
val bg = ArrayList<Event>(30_000 * SCALE + 300)
repeat(30_000 * SCALE) {
val author = hex64(1, it % 300)
bg.add(ev(author, base + (mix(it * 31L) and 0x7fffffff) % span, 1))
}
// Target author: 200 kind-1 notes + 50 kind-7 reactions.
repeat(200) { bg.add(ev(target, base + (mix(it * 131L) and 0x7fffffff) % span, 1)) }
repeat(50) { bg.add(ev(target, base + (mix(it * 61L) and 0x7fffffff) % span, 7)) }
val t0 = System.nanoTime()
store.transaction { bg.forEach { insert(it) } }
val insertMs = (System.nanoTime() - t0) / 1e6
println("─ FsDriverSelectionBenchmark: ${bg.size} events (scale=$SCALE), seed %.0f ms ─".format(insertMs))
// The planner's own pick — expected to choose the author
// tree over the ~30k-entry kind-1 tree.
val kindDriven = Filter(authors = listOf(target), kinds = listOf(1), limit = 50)
time(store, "planner (cost-based pick)") { store.query<Event>(kindDriven).size }
// Reference: author tree walked explicitly, kind checked
// by the caller — the pick should match or beat this.
time(store, "author-driver emulation") {
store
.query<Event>(Filter(authors = listOf(target), limit = 250))
.asSequence()
.filter { it.kind == 1 }
.take(50)
.count()
}
// Floor: author-only shape, the cheapest walk of the tree.
val authorOnly = Filter(authors = listOf(target), limit = 50)
time(store, "author-only floor") { store.query<Event>(authorOnly).size }
} finally {
store.close()
if (root.exists()) {
Files.walk(root).use { it.sorted(Comparator.reverseOrder()).forEach { p -> Files.deleteIfExists(p) } }
}
}
}
private inline fun time(
store: FsEventStore,
label: String,
run: () -> Int,
) {
repeat(3) { run() }
val runs = 10
var rows = 0
val start = System.nanoTime()
repeat(runs) { rows = run() }
val ms = (System.nanoTime() - start) / 1e6 / runs
println(" %-32s %8.2f ms (%d rows)".format(label, ms, rows))
}
}
@@ -70,11 +70,11 @@ object Scenarios {
} }
val topAuthors = notesByAuthor.entries.sortedWith(compareByDescending<Map.Entry<String, Int>> { it.value }.thenBy { it.key }).map { it.key } val topAuthors = notesByAuthor.entries.sortedWith(compareByDescending<Map.Entry<String, Int>> { it.value }.thenBy { it.key }).map { it.key }
val hottestThread = val hotNotes =
eTagRefs.entries eTagRefs.entries
.sortedWith(compareByDescending<Map.Entry<String, Int>> { it.value }.thenBy { it.key }) .sortedWith(compareByDescending<Map.Entry<String, Int>> { it.value }.thenBy { it.key })
.firstOrNull() .map { it.key }
?.key val hottestThread = hotNotes.firstOrNull()
val mostMentioned = val mostMentioned =
pTagRefs.entries pTagRefs.entries
.sortedWith(compareByDescending<Map.Entry<String, Int>> { it.value }.thenBy { it.key }) .sortedWith(compareByDescending<Map.Entry<String, Int>> { it.value }.thenBy { it.key })
@@ -86,6 +86,23 @@ object Scenarios {
.firstOrNull() .firstOrNull()
?.key ?.key
// The author that most often tags the most-mentioned pubkey — a
// conversation pair for the tag ∩ author (DM-room) query shape.
val conversationPeer =
mostMentioned?.let { me ->
val byAuthor = HashMap<String, Int>()
for (e in events) {
if (e.kind != 1 || e.pubKey == me) continue
if (e.tags.any { it.size >= 2 && it[0] == "p" && it[1] == me }) {
byAuthor.merge(e.pubKey, 1, Int::plus)
}
}
byAuthor.entries
.sortedWith(compareByDescending<Map.Entry<String, Int>> { it.value }.thenBy { it.key })
.firstOrNull()
?.key
}
// Evenly spread sample of note ids — a "fetch these 100 events" batch. // Evenly spread sample of note ids — a "fetch these 100 events" batch.
val idSample = val idSample =
if (noteIds.size <= 100) { if (noteIds.size <= 100) {
@@ -159,6 +176,33 @@ object Scenarios {
), ),
) )
} }
if (mostMentioned != null && conversationPeer != null) {
// The tag ∩ author ∩ kind shape (65 client assembler call
// sites: NIP-04 DM rooms, reports-by-follows, follows-scoped
// community feeds). Modeled on kind 1 because public corpora
// carry no DMs; the index path exercised is identical.
add(
Scenario(
"conversation",
"notes by one author tagging the most-mentioned pubkey (DM-room shape)",
Filter(kinds = listOf(1), authors = listOf(conversationPeer), tags = mapOf("p" to listOf(mostMentioned)), limit = 500),
),
)
}
if (hotNotes.size > 1) {
// Large-IN tag watcher: per-value streams come sorted off the
// tag index but their union does not, exposing whether the
// store collects+sorts or merges. 150 values stays inside
// strfry's default 200-element filter cap.
val watched = hotNotes.take(150)
add(
Scenario(
"reactions-watch",
"reactions on the ${watched.size} hottest notes (visible-feed reaction watcher)",
Filter(kinds = listOf(7), tags = mapOf("e" to watched), limit = 500),
),
)
}
topHashtag?.let { topHashtag?.let {
add( add(
Scenario( Scenario(