Verifies the FilterIndex path in LiveEventStore: each event matches exactly ONE of N subscribers, so without an index the relay walks all N subs per event (O(N)); with the author-keyed index the lookup is O(1) and per-event latency stays flat as N grows. Different from fanoutLatency (broadcast: 1 event reaches every sub). Here per-event work is lookup-bound rather than delivery-bound. Configurable via -DfanoutScalingSubs (comma list) and -DfanoutScalingEvents. Measured (laptop, JDK 21, full sweep): 100 subs, 2k events: p50=1.35ms p99=5.71ms 1000 subs, 2k events: p50=1.05ms p99=3.05ms 5000 subs, 1k events: p50=1.01ms p99=2.96ms p50 ~1ms across N — the predicted O(1) scaling. Default events count adapts downward at high N to stay below WebSocketSessionPump's 8192-frame outbound cap (test-client read side is the bottleneck above ~5k subs, not the relay). Plan doc updated with the measured numbers in the "How to verify" section.
8.6 KiB
Live broadcast: indexed filter matching for fanout
Status (2026-05-07): Phase 1 (relay server) and Phase 2 (Amethyst client
LocalCache.observables) are implemented. Phase 3 (per-projection dispatch fromObservableEventStore.changes) is left as future work — see "What's next" below.
Problem
Every accepted EVENT runs through LiveEventStore.newEventStream
(quartz/nip01Core/relay/server/LiveEventStore.kt) — a
MutableSharedFlow<Event> that every active subscription collects.
Each subscriber's collector then calls:
if (filters.any { it.match(newEvent) }) onEach(newEvent)
That's O(N_subscribers × N_filters_per_sub) per published event.
With 5k connections × ~3 filters average that's 15k Filter.match
calls per EVENT — and each Filter.match itself walks kinds,
authors, tag prefixes, since/until, etc. At 2k EPS ingest that's
~30M comparisons/sec.
The same shape recurs in two more places:
LocalCache.observablesinamethyst/.../LocalCache.kt— aConcurrentHashMap<Observable, Observable>of feed observers.refreshNewNoteObserversiterates every observer for every accepted event; each observer'snew()runsfilter.match.EventStoreProjectionunderquartz/.../cache/projection/— each projection collects everyStoreChange.InsertfromObservableEventStore.changesand runs its own filter list.
All three follow the pattern "many filter-bearing observers, one
incoming event — find which observers match". Today that's a per-event
walk over N observers; with an inverted index it becomes a few hash
lookups followed by Filter.match only on the (small) candidate set.
Solution
FilterIndex<S>
quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/filters/FilterIndex.kt
is a generic, KMP-friendly inverted index parameterised by the
subscriber type. It lives next to Filter.kt, not under relay/server/,
because it isn't relay-specific.
API:
class FilterIndex<S : Any> {
fun register(filter: Filter, subscriber: S)
fun register(filters: List<Filter>, subscriber: S)
fun registerUnindexed(subscriber: S) // predicate-only callers
fun unregister(subscriber: S)
fun candidatesFor(event: Event): Set<S> // index lookup
fun forEach(action: (S) -> Unit) // full iteration (delete paths)
fun size(): Int
fun isEmpty(): Boolean
}
State is held in a single AtomicReference<State<S>> (BanStore-style)
so reads in the hot path are wait-free. Writes copy-on-write; that's
fine because writes are subscription-rate (rare) while reads are
event-rate (frequent).
Indexing strategy: each filter contributes entries to one
dimension — the most selective indexable field. Picking one dimension
instead of all of them avoids over-counting subscribers in
candidatesFor and minimises bucket churn:
ids(most selective; an id matches one event).authors.- The first single-letter tag in
tags(thentagsAll). kinds.- None of the above → registered into the unindexed pool.
Multi-filter registrations OR the per-filter selections together,
which mirrors filters.any { it.match(...) }. Negative constraints
(since / until / tagsAll / saturated limit) stay in
Filter.match — the index produces a super-set, callers post-filter.
Phase 1: LiveEventStore
LiveEventStore.kt no longer uses a MutableSharedFlow. Instead:
- One
FilterIndex<LiveSubscription>shared across all REQs. insert()callsindex.candidatesFor(event)then runsFilter.matchon each candidate. Synchronous delivery — callers (the relay'sRelaySession) keep thedelivercallback cheap (queue-to-outbound).query()registers the subscription into the index before the historical replay starts (closes the same race the previousonSubscriptionhandoff closed), runs the replay, signals EOSE, thenawaitCancellation(). Live events arrive via index dispatch during the suspend;finally { index.unregister(sub) }cleans up.- Dedupe set during the historical phase is held in an
AtomicReference<HashSet<String>?>so the live-dispatch coroutine sees the post-EOSE handoff promptly.
Phase 2: LocalCache.observables
amethyst/.../LocalCache.kt swapped its
ConcurrentHashMap<Observable, Observable> for a
FilterIndex<Observable>. observeNotes / observeEvents /
observeNewEvents(filter: Filter) now register(filter, observer);
the predicate-only observeNewEvents(predicate) overload uses
registerUnindexed because the index can't introspect an opaque
predicate. Dispatch:
refreshNewNoteObserversiteratesobservables.candidatesFor(event)instead of every observer.refreshDeletedNoteObserversstill usesobservables.forEach { ... }— the index doesn't help on the delete path because every observer might hold the deleted note in its result set, and there's no event-shape to consult.
How to verify
geode.perf.LoadBenchmark.fanoutScaling (added):
- N connections, each subscribes to
{authors: [pk_i], kinds: [1]}. - Publish events round-robin across the N pubkeys; each event matches exactly one subscriber.
- Measure end-to-end latency p50/p99 for N ∈ {100, 1000, 5000}.
The benchmark adapts the per-N event count downward to stay below
geode's WebSocketSessionPump.MAX_OUTGOING_BUFFER (8192 frames per
session). The single-WS test subClient can't drain (subs + events)
frames at full firehose rate above ~5k subs in a shared test JVM —
that's a test-infra ceiling, not a relay one. Override with
-DfanoutScalingEvents=N to push past the default.
Measured numbers from a development laptop (Linux, JDK 21, full sweep):
| N subs | events | mean (ms) | p50 (ms) | p99 (ms) | p999 (ms) |
|---|---|---|---|---|---|
| 100 | 2 000 | 1.63 | 1.35 | 5.71 | 20.49 |
| 1 000 | 2 000 | 1.12 | 1.05 | 3.05 | 9.40 |
| 5 000 | 1 000 | 1.07 | 1.01 | 2.96 | 7.38 |
p50 stays at ~1 ms across all three N values — the index is producing the predicted O(1)-per-event scaling. The minor p99 spread comes from GC pauses and OkHttp scheduler jitter on the test client, not from linear per-event work in the relay.
For the LocalCache side, a similar benchmark would publish events matching one of M observers and measure dispatch cost as M grows. Not yet added — Phase 2 candidate for a follow-up.
Risks
- Subscription churn: re-subscribing on every page (the way some
client features work) means many index insert/remove operations.
COW on a single
AtomicReferencemakes each write a full inner-map copy; benchmark this path on a busy account to confirm the constants stay reasonable. - Tag explosion: an EVENT with many
e/ptags hits many tag buckets. The single-dimension-per-filter selection caps how many buckets contribute candidates per event — registering on the most selective dimension means tag-keyed filters typically pick one specific tag value, so an event's tag walk only finds filters registered under that exact(letter, value). - Memory: the index is a per-bucket set of subscription handles. At 5k subs × average 1 dimension × a handful of values per filter, ~10k–20k entries — negligible.
- Correctness fence:
registerhappens before historical replay on the relay side, and inside thecallbackFlow'sregister/awaitClose { unregister }pair on the client side. Events arriving mid-historical are deduped viaseenIds(relay) or are a non-issue because the clientobserve*flow seeds viainit()before any new event can fire. - Filters with no narrowing field (e.g.
{since: X}) fall into the unindexed pool and behave like today — every event reaches them. That's the worst case; it's not worse than the pre-index baseline. - AddressableEvent / replaceable v2 path: an observer holding v1
whose filter doesn't match v2 won't be in
candidatesFor(v2). Today such an observer wouldn't update its membership either (filter.match(v2) == falseshort-circuits before the re-emit branch). Pre-existing behaviour preserved.
What's next (Phase 3)
ObservableEventStore.changes is still a SharedFlow<StoreChange>
that every projection collects. To use the index there, the dispatcher
between _changes.emit and the per-projection collectors would consult
a FilterIndex<EventStoreProjection<*>>-style index and only deliver
to interested projections. Doable, but the SharedFlow contract is
public; replacing it is a larger refactor than Phase 1/2 and the ROI
is lower (per-projection apply is already small). Treat as a follow-up.