perf: cache the sealed negentropy snapshot across NEG-OPENs

A NIP-77 server session rebuilt its reconciliation structure from
scratch on every NEG-OPEN: full id+created_at scan, per-entry hex
decode into a fresh StorageVector, O(n log n) seal. That cost grows
with the corpus and is paid even when nothing changed — the exact
shape of a periodic mirror's heartbeat, where N peers reconcile the
same broad filter over and over. relayBench measured 342 ms per
identical-set reconcile at 50k events vs strfry's 26 ms off its
always-current LMDB tree.

Reconciliation only *reads* the sealed storage, so one instance can
back any number of concurrent sessions:

- NegentropyServerSession now accepts a pre-sealed IStorage (the
  List<IdAndTime> constructor remains and delegates).
- SessionBackend.sealedNegentropyStorage() builds + seals (null when
  the set exceeds maxSyncEvents); LiveEventStore overrides it with a
  single-slot cache keyed by (filter set, write generation) plus a 30s
  TTL. The generation bumps on every accepted ingest; the TTL bounds
  staleness from delete paths the counter can't see (expiration
  sweeps, admin purges) — negentropy snapshots are point-in-time sets,
  so seconds of staleness only means a peer briefly re-offers ids.
- NegSessionRegistry.open consumes the shared sealed storage;
  over-cap NEG-ERR behavior unchanged (strfry parity).

relayBench gains a 'heartbeat' measurement — the identical-set
reconcile repeated immediately with no writes in between. At 50k
events: geode 342 ms -> 27.8 ms vs strfry 21.1 ms (near parity; was
13x). Cold reconciles (first open after a write) are unchanged.

Also fixes the GeodeVsStrfryNegentropySyncTest fixture to write
'nofiles = 0' so the opt-in interop test can boot strfry inside
containers with a low RLIMIT_NOFILE hard cap; the interop suite passes
against strfry v1-b80cda3 with the cache in place.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NeoCvXnTxsKzqurkmjdC46
This commit is contained in:
Claude
2026-07-03 18:37:53 +00:00
parent 3c6d36cde1
commit f2174bdab3
7 changed files with 152 additions and 14 deletions
@@ -99,6 +99,7 @@ class GeodeVsStrfryNegentropySyncTest {
relay {
bind = "127.0.0.1"
port = $port
nofiles = 0
}
""".trimIndent(),
)
@@ -98,9 +98,12 @@ class NegSessionRegistry(
// NIP-77: same-subId OPEN replaces any prior session.
sessions.remove(cmd.subId)
val cap = settings.maxSyncEvents
val entries = store.snapshotIdsForNegentropy(filters, maxEntries = cap)
if (entries.size > cap) {
// Sealed storage may come from the backend's snapshot cache — a
// repeated NEG-OPEN of the same filter with no writes in between
// (the mirror-heartbeat pattern) skips the scan + sort entirely.
// `null` = matching set exceeds the cap (strfry-parity error).
val sealedStorage = store.sealedNegentropyStorage(filters, maxEntries = settings.maxSyncEvents)
if (sealedStorage == null) {
send(NegErrMessage(cmd.subId, "blocked: too many query results"))
return
}
@@ -108,7 +111,7 @@ class NegSessionRegistry(
val session =
NegentropyServerSession(
subId = cmd.subId,
localEntries = entries,
sealedStorage = sealedStorage,
frameSizeLimit = settings.frameSizeLimit,
)
sessions[cmd.subId] = session
@@ -20,6 +20,7 @@
*/
package com.vitorpamplona.quartz.nip01Core.relay.server.backend
import com.vitorpamplona.negentropy.storage.IStorage
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.relay.filters.FilterIndex
@@ -27,9 +28,12 @@ import com.vitorpamplona.quartz.nip01Core.store.IEventStore
import com.vitorpamplona.quartz.nip01Core.store.IdAndTime
import com.vitorpamplona.quartz.nip01Core.store.RawEvent
import com.vitorpamplona.quartz.nip50Search.strippingSearchExtensions
import com.vitorpamplona.quartz.utils.TimeUtils
import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.awaitCancellation
import kotlin.concurrent.atomics.AtomicBoolean
import kotlin.concurrent.atomics.AtomicLong
import kotlin.concurrent.atomics.AtomicReference
import kotlin.concurrent.atomics.ExperimentalAtomicApi
/**
@@ -90,6 +94,7 @@ class LiveEventStore(
) {
ingest.submit(event) { outcome ->
if (outcome is IEventStore.InsertOutcome.Accepted) {
writeGeneration.addAndFetch(1L)
fanout(event)
}
onComplete(outcome)
@@ -303,5 +308,72 @@ class LiveEventStore(
override suspend fun snapshotIdsForNegentropy(
filters: List<Filter>,
maxEntries: Int?,
<<<<<<< HEAD
): List<IdAndTime> = store.snapshotIdsForNegentropy(filters.strippingSearchExtensions(), maxEntries)
=======
): List<IdAndTime> = store.snapshotIdsForNegentropy(filters, maxEntries)
// ------------------------------------------------------------------
// NIP-77 snapshot cache
// ------------------------------------------------------------------
/**
* Bumped after every accepted write. A cached negentropy snapshot is
* only valid while this hasn't moved. Deletion paths that bypass the
* ingest queue (expiration sweeps, NIP-86 admin purges) don't bump it,
* which is why cache entries also carry a short TTL: a snapshot is a
* point-in-time set by NIP-77's nature, and a few seconds of staleness
* only means a peer momentarily re-offers ids the relay just dropped.
*/
private val writeGeneration = AtomicLong(0L)
private class CachedSnapshot(
val filterKey: String,
val generation: Long,
val builtAt: Long,
val storage: IStorage?,
)
private val snapshotCache = AtomicReference<CachedSnapshot?>(null)
/**
* Serves repeated NEG-OPENs of the same filter from one sealed
* storage as long as no write landed in between (single slot — the
* mirror-heartbeat pattern is many peers reconciling the same broad
* filter, not many filters). Rebuilding on every open costs a full
* scan + O(n log n) seal that grows with the corpus: relayBench
* measured 342 ms per identical-set reconcile at 50k events vs
* strfry's 26 ms off its always-current tree.
*/
override suspend fun sealedNegentropyStorage(
filters: List<Filter>,
maxEntries: Int,
): IStorage? {
val generation = writeGeneration.load()
val key = filters.joinToString("") { it.toJson() } + "cap=$maxEntries"
val now = TimeUtils.now()
val cached = snapshotCache.load()
if (cached != null &&
cached.filterKey == key &&
cached.generation == generation &&
now - cached.builtAt <= SNAPSHOT_TTL_SECONDS
) {
return cached.storage
}
val built = super.sealedNegentropyStorage(filters, maxEntries)
snapshotCache.store(CachedSnapshot(key, generation, now, built))
return built
}
private companion object {
/**
* Ceiling on how long a cached snapshot may serve NEG-OPENs even
* with no observed writes — bounds staleness from delete paths
* the generation counter can't see.
*/
const val SNAPSHOT_TTL_SECONDS = 30L
}
>>>>>>> 55139747 (perf: cache the sealed negentropy snapshot across NEG-OPENs)
}
@@ -20,12 +20,14 @@
*/
package com.vitorpamplona.quartz.nip01Core.relay.server.backend
import com.vitorpamplona.negentropy.storage.IStorage
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.CountResult
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.store.IEventStore
import com.vitorpamplona.quartz.nip01Core.store.IdAndTime
import com.vitorpamplona.quartz.nip01Core.store.RawEvent
import com.vitorpamplona.quartz.nip77Negentropy.NegentropyServerSession
/**
* The data plane a [RelaySession] talks to: how REQ/COUNT are answered, how
@@ -118,4 +120,25 @@ interface SessionBackend {
filters: List<Filter>,
maxEntries: Int?,
): List<IdAndTime> = emptyList()
/**
* NIP-77 snapshot as a **sealed** negentropy storage, ready to back a
* server session. Returns `null` when the matching set exceeds
* [maxEntries] (the caller answers NEG-ERR, strfry-parity).
*
* Reconciliation only reads the storage, so implementations may hand
* the same sealed instance to any number of concurrent sessions —
* [LiveEventStore] caches the most recent one and reuses it until a
* write invalidates it, which turns the periodic-mirror heartbeat
* ("anything new since last sync?") from a full scan + sort per
* NEG-OPEN into a cache hit. The default builds fresh per call.
*/
suspend fun sealedNegentropyStorage(
filters: List<Filter>,
maxEntries: Int,
): IStorage? {
val entries = snapshotIdsForNegentropy(filters, maxEntries)
if (entries.size > maxEntries) return null
return NegentropyServerSession.sealVector(entries)
}
}
@@ -21,6 +21,7 @@
package com.vitorpamplona.quartz.nip77Negentropy
import com.vitorpamplona.negentropy.Negentropy
import com.vitorpamplona.negentropy.storage.IStorage
import com.vitorpamplona.negentropy.storage.StorageVector
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.store.IdAndTime
@@ -50,19 +51,22 @@ import com.vitorpamplona.quartz.utils.Hex
*/
class NegentropyServerSession(
val subId: String,
localEntries: List<IdAndTime>,
/**
* A **sealed** storage. Reconciliation only reads it (range
* fingerprints), so one sealed instance may safely back any number
* of concurrent sessions — that is what lets the relay cache the
* snapshot across NEG-OPENs instead of rebuilding per session.
*/
sealedStorage: IStorage,
frameSizeLimit: Long = DEFAULT_FRAME_SIZE_LIMIT,
) {
private val storage = StorageVector()
private val negentropy: Negentropy
constructor(
subId: String,
localEntries: List<IdAndTime>,
frameSizeLimit: Long = DEFAULT_FRAME_SIZE_LIMIT,
) : this(subId, sealVector(localEntries), frameSizeLimit)
init {
for (entry in localEntries) {
storage.insert(entry.createdAt, entry.id)
}
storage.seal()
negentropy = Negentropy(storage, frameSizeLimit)
}
private val negentropy: Negentropy = Negentropy(sealedStorage, frameSizeLimit)
companion object {
/**
@@ -72,6 +76,16 @@ class NegentropyServerSession(
*/
const val DEFAULT_FRAME_SIZE_LIMIT: Long = 500_000L
/** Builds and seals a [StorageVector] from `(created_at, id)` pairs. */
fun sealVector(localEntries: List<IdAndTime>): StorageVector {
val storage = StorageVector()
for (entry in localEntries) {
storage.insert(entry.createdAt, entry.id)
}
storage.seal()
return storage
}
/**
* Convenience for callers that hold full [Event] objects
* (mostly tests + relay-relay sync paths). Production server
@@ -239,6 +239,15 @@ object Report {
higherIsBetter = false,
detail = pair.identicalReconcile.mapValues { (_, s) -> "${s.rounds} rounds, ${bytesHuman(s.wireBytes)}" },
)
if (pair.repeatReconcile.isNotEmpty()) {
appendLine(dim(" heartbeat: same reconcile again, no writes in between:"))
metricRows(
pair.repeatReconcile.map { (name, s) -> name to s.ms },
"ms",
higherIsBetter = false,
detail = pair.repeatReconcile.mapValues { (_, s) -> "${s.rounds} rounds, ${bytesHuman(s.wireBytes)}" },
)
}
}
}
}
@@ -79,6 +79,14 @@ object SyncBenchmark {
val transferredToA: Int,
val transferredToB: Int,
val identicalReconcile: Map<String, ReconcileStats>,
/**
* A second identical-set reconcile fired immediately after the
* first with no writes in between — the periodic-mirror
* heartbeat ("anything new?" × N peers). Relays that keep a
* live/cached reconciliation structure answer from it here;
* relays that rebuild per NEG-OPEN pay the full cost again.
*/
val repeatReconcile: Map<String, ReconcileStats> = emptyMap(),
val converged: Boolean,
val error: String? = null,
)
@@ -255,6 +263,13 @@ object SyncBenchmark {
log(" re-reconciling identical sets…")
val identicalA = reconcile(urlA, effective, http)
val identicalB = reconcile(urlB, effective, http)
// 4. Heartbeat: same reconcile again, immediately, no writes in
// between — measures whether the server keeps its reconciliation
// structure warm across NEG-OPENs.
log(" repeating identical-set reconcile (warm)…")
val repeatA = reconcile(urlA, effective, http)
val repeatB = reconcile(urlB, effective, http)
val converged =
diffCorrect &&
identicalA.error == null &&
@@ -275,6 +290,7 @@ object SyncBenchmark {
transferredToA = toA.size,
transferredToB = toB.size,
identicalReconcile = mapOf(nameA to stats(identicalA), nameB to stats(identicalB)),
repeatReconcile = mapOf(nameA to stats(repeatA), nameB to stats(repeatB)),
converged = converged,
)
}