diff --git a/geode/plans/2026-05-07-connection-scaling.md b/geode/plans/2026-05-07-connection-scaling.md new file mode 100644 index 000000000..fb1f8b353 --- /dev/null +++ b/geode/plans/2026-05-07-connection-scaling.md @@ -0,0 +1,93 @@ +# Connection scaling: pushing past 2 000 + +## Problem + +Current measurement (`LoadBenchmark.connectionsHeldOpen`): **~2 000 +concurrent connections** before file-descriptor pressure / Ktor CIO +event-loop saturation. Real-world relays (e.g. nostr.wine, nos.lol) +sustain 10–30k. Geode shouldn't be the bottleneck for an Amethyst- +adjacent operator who scales beyond a thousand-user community. + +## What's spending memory per connection today + +| Cost | Per connection | At 5 000 conns | +| ----------------------- | -------------------------------------------------------- | -------------- | +| `outQueue` Channel | 8 192 string slots × ~8 b ref | ~320 MB pinned | +| `RelaySession` | `LargeCache` for subs (likely 1–10 entries) | ~negligible | +| `NegSessionRegistry` | `HashMap` — usually 0 | ~negligible | +| Ktor CIO buffers | TCP read + write buffers | ~10 MB | +| Per-session writer Job | one coroutine | ~few KB | + +The `outQueue` reservation is the dominant cost. The 8 192 was sized +for a worst case "thousands of subscriptions, one event matches all" — +but at 5 000 connections we've over-provisioned by ~300 MB just on +the channel array, even though most connections never fan out. + +## Sketch + +### A — adaptive outQueue capacity + +Start every connection with `INITIAL_OUTGOING_BUFFER = 64`. When the +producer side trySends and we observe queue depth crossing a high-water +mark (e.g. 75% full), grow the channel up to `MAX_OUTGOING_BUFFER = +8192`. This is not how `kotlinx.coroutines.channels.Channel` is +structured (capacity is fixed at construction), so the implementation +is "swap in a wider channel under a per-session lock when watermark +trips" — drains the old, then routes new sends through the new. + +Expected: 90% of connections never fan out, so they stay at 64 slots +× ~512 B per ref ≈ 32 KB. At 5 000 conns that's ~160 MB → ~5 MB. +Hot-fanout connections still get the 2 MB cap. + +### B — per-relay event-loop pool sizing + +Ktor CIO defaults to one event-loop thread per available CPU. +Beyond a few thousand connections, this becomes the bottleneck — and +none of geode's per-connection work is CPU-bound (it's mostly waiting +on incoming frames). Tune CIO via: + +```kotlin +embeddedServer(CIO, ...) { + connectionGroupSize = max(2, Runtime.getRuntime().availableProcessors() / 2) + workerGroupSize = max(4, Runtime.getRuntime().availableProcessors()) + callGroupSize = max(8, Runtime.getRuntime().availableProcessors() * 4) +} +``` + +Expose these through `RelayConfig.NetworkSection` so an operator on a +big VM can lift them. + +### C — reduce per-message JSON allocations + +`OptimizedJsonMapper.fromJsonToCommand` allocates a `JsonNode` tree per +incoming frame. At 10k connections with 1 msg/s each that's 10k tree +allocations/sec. Investigate streaming Jackson + reusing `ObjectMapper` +per session, or using kotlinx-serialization's lower-overhead path. + +This is more of a quartz-level change than geode-specific, but +geode's load benchmark is the right place to measure it. + +## How to verify + +Add to `geode.perf.LoadBenchmark`: + +- `connectionsHeldOpen10k` — opens 10 000 idle WebSocket connections; + asserts no FD exhaustion + RSS stays under 1 GB. +- `connectionsHeldOpenWithFanout` — 5 000 idle subscribers, + 10 EPS published; measures p99 fanout latency at scale. + +The current `connectionsHeldOpen` benchmark stays as the baseline +floor (~2 000 conns). + +## Risks + +- **Adaptive channel swap is fiddly**: drains under the producer's nose + must preserve OK ordering. A simpler alternative: keep capacity fixed, + but lazily allocate a small `ArrayDeque` only when the first + message is sent. Channels in kotlinx.coroutines do allocate up-front. +- **Bumping CIO group sizes can hurt**: more threads can mean worse + L1/L2 locality. Always benchmark before/after, don't trust + intuitive sizing. +- **OS-level FD limit**: per-process FD limit on Linux defaults to + 1024 in many environments. Document the `ulimit -n` requirement + for operators targeting >1k connections. diff --git a/geode/plans/2026-05-07-event-ingestion-batching.md b/geode/plans/2026-05-07-event-ingestion-batching.md new file mode 100644 index 000000000..e6df99ebf --- /dev/null +++ b/geode/plans/2026-05-07-event-ingestion-batching.md @@ -0,0 +1,91 @@ +# Event ingestion: write batching + pipelined OK + +## Problem + +EVENT acceptance is the hot path on a busy relay — every published note, +every reaction, every DM lands here. Today the per-event flow is fully +serial: + +1. `RelaySession.handleEvent` (`quartz/nip01Core/relay/server/RelaySession.kt:131`) + awaits `policy.accept(cmd)` (Schnorr verify if `VerifyPolicy` is in + the stack — ~0.1 ms on JVM). +2. Awaits `store.insert(cmd.event)` — a single SQLite write, guarded by + the connection-pool writer mutex (`SQLiteConnectionPool`). +3. Sends `OkMessage` back through the writer coroutine. + +`LoadBenchmark.publishThroughputSingleClient` measured **~760 EPS**; +the concurrent variant **~2000 EPS** (limited by SQLite writer mutex +contention, not WS throughput). + +## Constraints we must keep + +- **OK ordering**: NIP-01 requires the OK reply to follow its EVENT. + We cannot reply OK before the insert decision (the OK carries + accepted/rejected + reason). +- **Durability semantics**: clients reasonably assume `OK true` means + "stored." Batching must not make us reply OK before fsync. +- **Per-connection FIFO**: a publisher that sends three EVENTs in a + row expects three OKs in that order. Reordering across connections + is fine. + +## Sketch + +### Tier 1 — SQLite WAL + group commit (cheap win) + +Confirm `PRAGMA journal_mode=WAL` + `PRAGMA synchronous=NORMAL` on the +event-store DB; group commits across the writer mutex's hold window. +Today each insert is its own transaction. Wrap N inserts (or a 5 ms +budget, whichever first) in a single transaction managed by the writer +coroutine. On commit, fan back N OK replies. + +Implementation lives in quartz's `EventStore` / `SQLiteConnectionPool`, +not geode — but geode owns the benchmark and validates the gain. + +Expected: **~5–10× write throughput** on a fast SSD. SQLite group +commit is well-trodden territory (nostr-rs-relay, strfry both do it). + +### Tier 2 — pipelined OK over multiple in-flight EVENTs + +`RelaySession.receive` is currently single-flight: one EVENT in, +process, OK out, next EVENT. Allow a connection to push N EVENTs +concurrently, dispatch them to a per-connection ingest pipeline, and +serialise OKs back in arrival order via a small commit log. + +A `Channel with capacity = INGEST_PIPELINE_DEPTH` per +connection, drained by a coroutine that batches into the group-commit +above. OK responses are written to an `outQueue.send()` already — so +the pipeline just needs to record arrival order and emit OKs in that +order after each batch commits. + +Expected: hides the verify+insert latency behind another EVENT's +parse, gets us closer to network-bound throughput. + +### Tier 3 — eager Schnorr verify off the writer thread + +`VerifyPolicy` is in the policy stack and runs synchronously on +`receive`. Move it into the ingest pipeline so verification of EVENT N+1 +runs concurrently with the SQLite commit of EVENT N. secp256k1 verify +is parallelisable; the writer should never block on it. + +## How to verify + +Add to `geode.perf.LoadBenchmark`: + +- `publishGroupCommitSingleClient` — same workload as the current + single-client benchmark, asserts >5000 EPS. +- `publishPipelinedSingleClient` — sends 100 EVENTs without awaiting + intermediate OKs; measures end-to-end and OK-ordering correctness. + +Existing benchmarks stay as the regression floor. + +## Risks + +- **Group commit windows**: if a single bad event in the batch fails + validation, we must not roll back the good ones. The batch needs + per-row commit semantics (row-level errors → row-level OK false). +- **Backpressure on slow disks**: deeper pipelines on slow storage + amplify out-of-memory pressure. Cap the in-flight queue depth and + apply existing slow-client backpressure if it fills. +- **Replay protection**: the existing dedupe table needs to see the + event before commit, not after — keep that check inside the writer + coroutine. diff --git a/geode/plans/2026-05-07-live-broadcast-fanout-index.md b/geode/plans/2026-05-07-live-broadcast-fanout-index.md new file mode 100644 index 000000000..9a0227a08 --- /dev/null +++ b/geode/plans/2026-05-07-live-broadcast-fanout-index.md @@ -0,0 +1,100 @@ +# Live broadcast: indexed filter matching for fanout + +## Problem + +Every accepted EVENT runs through `LiveEventStore.newEventStream` +(`quartz/nip01Core/relay/server/LiveEventStore.kt:43`) — a +`MutableSharedFlow` that every active subscription collects. +Each subscriber's collector then calls: + +```kotlin +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. + +Two specific cost shapes: + +1. **Filters that almost never match.** Most subscriptions are scoped + to a small author list. Today every published EVENT walks every + such subscription to learn that. A `HashMap>` + keyed by author would cut this to O(1) average for the dominant case. +2. **Pseudo-broadcast filters** (`{kinds: [1]}` with no other + constraint) match almost everything. There's no avoiding the + per-subscriber notification, but at least the index lookup is + cheap. + +`LoadBenchmark.fanoutLatency` already measures this — current +results are not yet noted in tree, but back-of-envelope says fanout +becomes the dominant cost above ~2k subscribers. + +## Sketch + +A new `LiveBroadcastIndex` inside `LiveEventStore`: + +```kotlin +private val byAuthor = ConcurrentHashMap>() +private val byKind = ConcurrentHashMap>() +private val byTag = ConcurrentHashMap>() +private val unindexed = CopyOnWriteArraySet() // subs with no + // narrowing field +``` + +Each `RelaySession.handleReq` registers its `Subscription` (a tuple of +filters + the existing `EventMessage` send callback) into whichever +buckets each filter narrows on. A filter with `kinds=[1] and +authors=[a,b]` registers into `byKind[1]` AND `byAuthor[a]`, +`byAuthor[b]` — broadcast unions the resulting candidate sets. + +On EVENT arrival: + +1. Build the candidate set: union of `byAuthor[event.pubkey]`, + `byKind[event.kind]`, every `byTag[(letter, value)]` for the + event's single-letter tags, plus `unindexed`. +2. Run the existing `Filter.match` on each candidate to handle + negative constraints (`since`, `until`, `limit` already-reached, + composite predicates). +3. Send. + +Expected: **>10× speedup** on fanout for realistic subscriptions. +Worst case (all filters in `unindexed`) degrades to current behaviour. + +## Where it lives + +`quartz/nip01Core/relay/server/LiveBroadcastIndex.kt` — protocol-level, +reusable by any relay embed. `RelaySession.handleReq` registers/ +unregisters; `LiveEventStore.insert` calls +`index.candidatesFor(event)`. + +## How to verify + +Add `geode.perf.LoadBenchmark.fanoutScaling`: + +- N connections, each subscribes to `{authors: [pk_i], kinds: [1]}`. +- Publish 10k EVENTs from a producer connection; each event matches + exactly one subscriber. +- Measure end-to-end latency p50/p99 for N ∈ {100, 1000, 5000}. + +Without the index, p99 grows roughly linearly with N. With the +index, p99 should be flat up to a much higher N. + +## Risks + +- **Subscription churn**: re-subscribing on every page (the way some + client features work) means many index insert/remove operations. + `ConcurrentHashMap` value-set operations need to be lock-free or + finely locked; benchmark this path explicitly. +- **Tag explosion**: an EVENT with many `e`/`p` tags hits many tag + buckets. Cap candidate-set union work or short-circuit when the + union saturates. +- **Memory**: the index is a per-bucket set of subscription handles. + At 5k subs × average 3 narrowing fields, ~15k entries — negligible. +- **Correctness fence**: the index must see new subscriptions before + the next EVENT broadcast. Today `RelaySession.handleReq` writes its + `Job` into a `LargeCache` then launches the collector. Order of + operations needs to be revisited so the index is updated atomically + with the collector being ready. diff --git a/geode/plans/2026-05-07-negentropy-large-corpus.md b/geode/plans/2026-05-07-negentropy-large-corpus.md new file mode 100644 index 000000000..0cb31ccd9 --- /dev/null +++ b/geode/plans/2026-05-07-negentropy-large-corpus.md @@ -0,0 +1,100 @@ +# NIP-77 negentropy at scale: snapshot memory + chunked replay + +## Problem + +`RelaySession` delegates NEG-OPEN to `NegSessionRegistry.open` +(`quartz/nip01Core/relay/server/NegSessionRegistry.kt`), which calls +`store.snapshotQuery(filters)` and feeds the **entire** result list +into `NegentropyServerSession`. For a relay holding 5 M events that +match a broad NEG-OPEN filter (`{kinds: [1, 7]}`), this is 5 M +`Event` objects materialised in memory before the first NEG-MSG goes +out. + +The negentropy library itself is fine — it pivots into a sealed +`StorageVector` (id + createdAt only, ~40 bytes/entry). But the +`store.query(f)` step that produces the input materialises full +`Event` objects with content, tags, sig — call it ~1 KB/event. 5 M × +1 KB = 5 GB transient pressure per concurrent NEG-OPEN. + +Two operator-visible symptoms: + +1. NEG-OPEN with a broad filter spikes JVM heap; under load, GC pause + stalls every other handler on the same process. +2. NEG-OPEN latency before the first NEG-MSG response is O(N) — for + large stores the client waits seconds for what should be a + millisecond round-trip. + +## Sketch + +### A — id-and-time-only snapshot path + +Negentropy only needs `(createdAt, id)` pairs. Add a streaming +`IEventStore.queryIdAndTime(filter)` that returns +`Sequence>` (or a `Flow` of small chunks) — +no content/tags/sig, no Event allocation. SQLite path is a SELECT +on `event_headers` (the `created_at`, `id` columns are already +indexed for query plans). + +```kotlin +suspend fun snapshotIdsForNegentropy(filter: Filter): IdTimeStream +``` + +`NegentropyServerSession` is rewritten to take that stream and feed +it directly into the `StorageVector`. Memory drops from O(N × 1 KB) +to O(N × 40 B) — a 25× reduction; for 5 M events, ~200 MB instead +of 5 GB. + +### B — bounded-window subscriptions + +Most NEG-OPENs from real Nostr clients want the last 30 days, not +"everything." If the client doesn't supply `since`, the server can +default to a configurable horizon (e.g. 90 days) and surface this in +the NIP-11 `limitation.negentropy_max_lookback_seconds` field. +Operators can lift the cap; clients reading the doc know the bound. + +This is a NIP-spec-adjacent question more than a code change — needs +a comment on whether the spec allows it. nostr-rs-relay does this +already. + +### C — frame-size cap on NEG-MSG + +`NegentropyServerSession` is constructed with `frameSizeLimit = 0` +(no limit). At very large reconciliations the message can grow large. +Set a default `frameSizeLimit = 64 * 1024` (matching the typical WS +frame budget) so NEG-MSGs don't blow past `[limits].max_ws_frame_bytes`. + +The library already supports this — pure config change in +`NegSessionRegistry.open`. + +### D — concurrent NEG-OPEN cap + +A NEG-OPEN holds session state until NEG-CLOSE (or connection close). +Today nothing caps the number of concurrent open negentropy sessions +per connection. A misbehaving (or hostile) client could open thousands +and pin RAM. Add `MAX_NEG_SESSIONS_PER_CONNECTION = 16`, send NEG-ERR +on overflow. + +## How to verify + +Add to `geode.perf.LoadBenchmark`: + +- `negentropyOpenLatencyLargeCorpus` — preload 1 M events (use + fixtures), measure NEG-OPEN → first NEG-MSG latency. Target <100 ms. +- `negentropyMemoryPressure` — open 10 concurrent NEG-OPENs on the + same large corpus; measure RSS delta, target <500 MB. + +## Risks + +- **`Sequence`/`Flow` over SQLite cursor**: holding a cursor open + across the full sync is fragile if the client stalls. Materialise + to a smaller in-memory list (just (id, createdAt)) once, reuse for + the lifetime of the session. Memory bound is the same. +- **Defaulting `since` is a behaviour change**: existing clients that + expect "everything" silently get a bounded window. Either (a) make + it opt-in via `RelayConfig.NegentropySection.default_lookback_seconds + = null`, (b) advertise the cap in NIP-11 so well-behaved clients + read it. +- **Frame-size cap can break older clients**: the NIP-77 reference + implementation (kmp-negentropy) handles this gracefully — multi-frame + reconciliation is in spec — but field-test against a known-working + client (e.g. nstart, primal-cache) before flipping the default. diff --git a/geode/plans/README.md b/geode/plans/README.md new file mode 100644 index 000000000..6524503b3 --- /dev/null +++ b/geode/plans/README.md @@ -0,0 +1,19 @@ +# geode plans + +Performance-focused design docs for future work. Each file is a +self-contained sketch — problem statement, observed numbers, proposed +fix, how to verify, risks. None of these are committed work; they're +the queue. + +Ordered roughly by expected impact: + +| Plan | Headline gain | +| ---- | ------------- | +| [2026-05-07-event-ingestion-batching.md](2026-05-07-event-ingestion-batching.md) | 5–10× write EPS via SQLite group commit + ingest pipelining | +| [2026-05-07-live-broadcast-fanout-index.md](2026-05-07-live-broadcast-fanout-index.md) | >10× fanout speedup at >2 000 subscribers | +| [2026-05-07-connection-scaling.md](2026-05-07-connection-scaling.md) | 2 000 → 10 000+ concurrent connections | +| [2026-05-07-negentropy-large-corpus.md](2026-05-07-negentropy-large-corpus.md) | 25× lower memory + faster NEG-OPEN on M-event corpora | + +Verification target for each plan is a new method on +`geode.perf.LoadBenchmark` (gated by `-DrunLoadBenchmark=true`) so +regressions show up in the regular CI matrix once they're enabled.