diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt index 2ff1e4c866..9c6acab5f1 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt @@ -260,7 +260,17 @@ fun main(args: Array) { if (upstreams.isEmpty()) { null } else { - MirrorWorker(upstreams, relay.server).also { it.start() } + MirrorWorker( + upstreams = upstreams, + server = relay.server, + // The store lets the catch-up reconcile against what we already + // hold (download only the diff, like `strfry sync`). + store = store, + // Production mirrors how strfry does it: NIP-77 "sync" catch-up + // for the historical window, then live REQ tail. Auto-falls back + // to paged REQ for upstreams without NIP-77. + negentropyBackfill = true, + ).also { it.start() } } // Periodic query-planner statistics refresh (`PRAGMA optimize`). diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt index eae9d20b0c..4e3e6c842d 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt @@ -23,6 +23,7 @@ package com.vitorpamplona.geode.mirror import com.vitorpamplona.quartz.nip01Core.core.Event import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient +import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.negentropySyncOrFetch import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.EventMessage import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.ReqCmd @@ -41,6 +42,7 @@ import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.cancel import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.channels.trySendBlocking import kotlinx.coroutines.delay import kotlinx.coroutines.launch import okhttp3.OkHttpClient @@ -132,11 +134,30 @@ class MirrorUpstream( class MirrorWorker( private val upstreams: List, private val server: NostrServer, + /** + * The local event set — read (never written) to enumerate the ids we + * already hold so the negentropy catch-up ([MirrorUpstream.backfillSeconds] + * > 0) reconciles against them and downloads only the diff, exactly like + * `strfry sync`. `null` reconciles against an empty local set (a full + * re-download of the window, which the store's unique-id constraint dedups) + * — used by tests that don't wire a store. + */ + private val store: IEventStore? = null, /** * Transport override for tests (e.g. `InProcessRelays`). Defaults to * a real OkHttp WebSocket per upstream. */ websocketBuilder: WebsocketBuilder? = null, + /** + * Whether the down catch-up uses NIP-77 negentropy (`strfry sync`) for the + * historical window before the live REQ tail takes over. The geode binary + * turns this on (see `Main`); it defaults **off** so the many existing + * MirrorWorker tests keep exercising the pure live-REQ path unchanged. + * When on, `negentropySyncOrFetch` automatically falls back to paged REQ + * against an upstream that doesn't speak NIP-77 — so "either mode" is + * transparent and needs no separate toggle. + */ + private val negentropyBackfill: Boolean = false, ) : AutoCloseable { private val scope = CoroutineScope(Dispatchers.IO + SupervisorJob()) @@ -308,12 +329,25 @@ class MirrorWorker( val scopedBase = (up.filter ?: Filter()).copy(since = null, limit = null) val initialSince = since - up.backfillSeconds + // strfry's two-phase model: a one-shot NIP-77 "sync" closes the + // historical [initialSince, now] gap (bounded, client-paced — it + // completes a bulk pull a plain REQ backfill can't), then the live + // REQ subscription tails everything new. With catch-up on, the down + // live sub starts at `now` (history is the sync's job); otherwise it + // replays from `initialSince` as before. The two windows overlap at + // `now`; the store's unique-id constraint dedups the seam. + val catchUp = negentropyBackfill && up.direction != MirrorDirection.UP && up.backfillSeconds > 0 + val downLiveSince = if (catchUp) since else initialSince + if (up.direction != MirrorDirection.UP) { - downSubs += startDown(i, up, scopedBase, initialSince, exchanged) + downSubs += startDown(i, up, scopedBase, downLiveSince, exchanged) } if (up.direction != MirrorDirection.DOWN) { startUp(up, scopedBase.copy(since = initialSince), exchanged) } + if (catchUp) { + scope.launch { runCatchUp(up, scopedBase, initialSince, since) } + } } client.connect() @@ -349,6 +383,88 @@ class MirrorWorker( } } + /** + * The NIP-77 "sync" phase — geode's equivalent of `strfry sync --dir down`, + * run once per down/both upstream to close the historical + * `[initialSince, until]` gap before (and alongside) the live REQ tail. + * + * Reconciles the local set against the upstream and downloads only the diff, + * **client-paced** so a fast upstream can't overrun the sink: a plain REQ + * backfill of a large set dies here (strfry kills a slow REQ client once its + * unsent-outbound buffer crosses `maxPendingOutboundBytes`), which is exactly + * why this uses negentropy. [INostrClient.negentropySyncOrFetch] falls back + * to paged REQ automatically when the upstream doesn't speak NIP-77, so the + * mirror is compatible with either kind of upstream with no config. + * + * A failure here is non-fatal: the live subscription keeps the mirror current + * and the reconnect watermark narrows any residual gap. + */ + private suspend fun runCatchUp( + up: MirrorUpstream, + scopedBase: Filter, + initialSince: Long, + until: Long, + ) { + val catchUpFilter = scopedBase.copy(since = initialSince, until = until) + // Reconcile against what we already hold in this window → download only + // the diff (like `strfry sync`). No store wired → empty local set → the + // whole window is downloaded and the store's unique-id constraint dedups. + val localEntries = store?.snapshotIdsForNegentropy(listOf(catchUpFilter)) ?: emptyList() + + // Bounded hand-off → one ingest consumer. `onEvent` can't suspend, so it + // blocks here when the sink falls behind; because negentropySyncOrFetch's + // own delivery pipeline is bounded, that backpressure reaches all the way + // to the upstream — no unbounded buffering (unlike the live-tail path). + val handoff = Channel(capacity = CATCHUP_HANDOFF) + val consumer = + scope.launch { + for (event in handoff) { + try { + server.ingest(event, up.trusted) { outcome -> + when (outcome) { + IEventStore.InsertOutcome.Accepted -> accepted.incrementAndGet() + is IEventStore.InsertOutcome.Rejected -> rejected.incrementAndGet() + } + } + } catch (e: CancellationException) { + throw e + } catch (e: Throwable) { + Log.w("MirrorWorker") { "catch-up ingest stopped for ${up.url.url}: ${e.message}" } + break + } + } + } + + try { + val result = + client.negentropySyncOrFetch( + relay = up.url, + filter = catchUpFilter, + localEntries = localEntries, + onEvent = { event -> + // Same containment as the live path: even a trusted + // upstream may only inject events inside the declared scope. + if (up.filter == null || up.filter.match(event)) { + handoff.trySendBlocking(event) + } else { + filtered.incrementAndGet() + } + }, + ) + Log.i("MirrorWorker") { + val how = if (result.pagedFallback) "paged REQ (upstream has no NIP-77)" else "negentropy" + "catch-up from ${up.url.url}: ${result.downloaded} events via $how" + } + } catch (e: CancellationException) { + throw e + } catch (e: Throwable) { + Log.w("MirrorWorker") { "catch-up from ${up.url.url} failed (live tail continues): ${e.message}" } + } finally { + handoff.close() + consumer.join() + } + } + /** Down direction: subscribe to the upstream, ingest what it sends. */ private fun startDown( index: Int, @@ -492,5 +608,13 @@ class MirrorWorker( * duplicate round trip. */ const val EXCHANGED_IDS_CAPACITY = 8_192 + + /** + * Depth of the catch-up hand-off between the negentropy download and the + * ingest consumer. Small: negentropySyncOrFetch is already internally + * backpressured, so this only smooths the seam — the bounded IngestQueue + * behind `server.ingest` is the real limiter. + */ + const val CATCHUP_HANDOFF = 4_096 } } diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorNegentropyCatchUpTest.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorNegentropyCatchUpTest.kt new file mode 100644 index 0000000000..1664e863a1 --- /dev/null +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorNegentropyCatchUpTest.kt @@ -0,0 +1,189 @@ +/* + * 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.geode.mirror + +import com.vitorpamplona.geode.KtorRelay +import com.vitorpamplona.geode.RelayEngine +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl +import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore +import com.vitorpamplona.quartz.utils.TimeUtils +import kotlinx.coroutines.delay +import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.withTimeoutOrNull +import kotlin.test.AfterTest +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertTrue + +/** + * Proves geode's two-phase mirror — the strfry model — over the real transport: + * a NIP-77 negentropy "sync" catch-up for the historical window, then the live + * REQ tail. + * + * The catch-up is isolated from the live tail by construction: the upstream is + * preloaded with **historical** events (created_at in the past) and the live + * subscription starts at `now`, so the live sub can never deliver them — only + * the catch-up covers `[now - backfill, now]`. Reaching the full count therefore + * proves the negentropy sync ran and completed. A fresh event published *after* + * boot then proves the live tail is still running alongside it. + */ +class MirrorNegentropyCatchUpTest { + private val upstreamStore = EventStore(null) + private val downstreamStore = EventStore(null) + + private val upstream = + RelayEngine(url = "ws://127.0.0.1:7894/".normalizeRelayUrl(), store = upstreamStore) + private val downstream = + RelayEngine(url = "ws://127.0.0.1:7895/".normalizeRelayUrl(), store = downstreamStore, parallelVerify = true) + + private var server: KtorRelay? = null + private var worker: MirrorWorker? = null + + @AfterTest + fun tearDown() { + worker?.close() + server?.stop(gracePeriodMillis = 0, timeoutMillis = 1_000) + upstream.close() + downstream.close() + } + + 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 suspend fun awaitCount( + target: Int, + timeoutMs: Long = 120_000, + ): Int { + var reached = 0 + var last = -1 + var stable = 0 + withTimeoutOrNull(timeoutMs) { + while (true) { + val c = downstreamStore.count(Filter()) + reached = c + if (c >= target) break + if (c == last) { + if (++stable >= 30) break // ~15s with no progress + } else { + stable = 0 + last = c + } + delay(500) + } + } + return reached + } + + @Test + fun negentropyCatchUpThenLiveTail() = + runBlocking { + val count = System.getProperty("catchUpN")?.toInt() ?: 3_000 + val now = TimeUtils.now() + val sig = "f".repeat(128) + + // Historical events: all safely in the past (but inside the backfill + // window), so the live-only sub (since = boot now) can NOT deliver + // them. Only the negentropy catch-up covers this range. + val history = + (0 until count).map { i -> + Event( + id = hex64(7, i), + pubKey = hex64(3, i % 500), + createdAt = now - 3_600 - (i % 1_000), + kind = 1, + tags = emptyArray(), + content = "h$i", + sig = sig, + ) + } + history.chunked(10_000).forEach { upstreamStore.batchInsert(it) } + assertEquals(count, upstreamStore.count(Filter()), "upstream preloaded") + + server = KtorRelay(upstream, host = "127.0.0.1", port = 7894).start() + + worker = + MirrorWorker( + upstreams = + listOf( + MirrorUpstream( + url = "ws://127.0.0.1:7894/".normalizeRelayUrl(), + trusted = true, + backfillSeconds = 86_400, + ), + ), + server = downstream.server, + store = downstreamStore, + negentropyBackfill = true, + ).also { it.start() } + + // Phase 1: the catch-up must deliver every historical event. + val afterCatchUp = awaitCount(count) + assertEquals(count, afterCatchUp, "negentropy catch-up dropped ${count - afterCatchUp} of $count historical events") + + // Phase 2: a fresh event published AFTER boot proves the live REQ + // tail is running alongside the catch-up. Published through the + // upstream server so it fans out to the live subscription. + val live = + Event( + id = hex64(9, 1), + pubKey = hex64(3, 0), + createdAt = TimeUtils.now(), + kind = 1, + tags = emptyArray(), + content = "live", + sig = sig, + ) + upstream.server.ingest(live, skipVerify = true) { } + + val afterLive = awaitCount(count + 1) + assertEquals(count + 1, afterLive, "live tail did not deliver the post-boot event") + assertTrue( + downstreamStore.count(Filter(ids = listOf(live.id))) == 1, + "the live event is present downstream", + ) + + println("─ MirrorNegentropyCatchUpTest: catch-up $count + live 1 = $afterLive delivered ─") + } +} diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt index 9d3a65eeee..930ad92d45 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientNegentropySyncExt.kt @@ -57,11 +57,12 @@ import kotlin.time.TimeSource /** * Outcome of a successful [negentropySync] run. * - * @property needCount ids the relay had that we lacked (i.e. everything that - * matched [Filter] on the relay — this sync always reconciles against an empty - * local set, so it downloads the full matched set). - * @property haveCount ids we had that the relay lacked. Always `0` here because - * the local set is empty; kept so the result mirrors a full NIP-77 reconcile. + * @property needCount ids the relay had that we lacked — the diff downloaded. + * With the default empty `localEntries` this is the relay's full matched set; + * pass the local set to reconcile incrementally and download only the diff. + * @property haveCount ids we had that the relay lacked. `0` unless `localEntries` + * is supplied (the downloader ignores this direction); kept so the result + * mirrors a full NIP-77 reconcile. * @property downloaded distinct events actually delivered through `onEvent`. * @property windows number of `created_at` windows the matched set was split * into (`1` when the relay reconciled the whole filter in one shot). @@ -78,9 +79,10 @@ class NegentropySyncResult( * one (deduped by id) through [onEvent]. A high-level wrapper over NIP-77 * negentropy that hides the parts that make the raw protocol painful to use: * - * 1. Reconciles the relay's matched set against an empty local set, **streaming** - * the ids the relay has straight into the download pipeline as each NIP-77 - * round arrives — the full id list is never materialised. + * 1. Reconciles the relay's matched set against [localEntries] (empty by default, + * which downloads the full matched set; pass the local ids to fetch only the + * diff), **streaming** the ids the relay has straight into the download + * pipeline as each NIP-77 round arrives — the full id list is never materialised. * 2. Downloads those ids through at most [maxConcurrentReqs] concurrent `REQ` * subscriptions of [fetchBatch] ids each, refilling as each `EOSE` arrives. The * reconciliation, the id queue and event delivery are all back-pressured, so a @@ -157,6 +159,7 @@ suspend fun INostrClient.negentropySync( idleTimeoutMs: Long = 120_000L, reconcileConcurrency: Int = 1, idBufferBatches: Int = maxConcurrentReqs * 4, + localEntries: List = emptyList(), onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, onEvent: (Event) -> Unit, ): NegentropySyncResult { @@ -190,6 +193,7 @@ suspend fun INostrClient.negentropySync( maxConcurrentReqs = maxConcurrentReqs, reconcileConcurrency = reconcileConcurrency, idBufferBatches = idBufferBatches, + localEntries = localEntries, onWindow = { windows.incrementAndFetch() }, // Only accumulate here; progress is reported from the // single consumer loop below so the user callback is never @@ -234,6 +238,7 @@ suspend fun INostrClient.negentropySync( idleTimeoutMs: Long = 120_000L, reconcileConcurrency: Int = 1, idBufferBatches: Int = maxConcurrentReqs * 4, + localEntries: List = emptyList(), onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, onEvent: (Event) -> Unit, ): NegentropySyncResult = @@ -246,6 +251,7 @@ suspend fun INostrClient.negentropySync( idleTimeoutMs = idleTimeoutMs, reconcileConcurrency = reconcileConcurrency, idBufferBatches = idBufferBatches, + localEntries = localEntries, onProgress = onProgress, onEvent = onEvent, ) @@ -291,6 +297,7 @@ suspend fun INostrClient.negentropySyncOrFetch( idleTimeoutMs: Long = 120_000L, reconcileConcurrency: Int = 1, idBufferBatches: Int = maxConcurrentReqs * 4, + localEntries: List = emptyList(), onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, onEvent: (Event) -> Unit, ): NegentropyOrFetchResult { @@ -319,6 +326,7 @@ suspend fun INostrClient.negentropySyncOrFetch( idleTimeoutMs = idleTimeoutMs, reconcileConcurrency = reconcileConcurrency, idBufferBatches = idBufferBatches, + localEntries = localEntries, onProgress = onProgress, ) { accept(it) } NegentropyOrFetchResult(delivered, pagedFallback = false, negentropy = result, fallbackCause = null) @@ -344,6 +352,7 @@ suspend fun INostrClient.negentropySyncOrFetch( idleTimeoutMs: Long = 120_000L, reconcileConcurrency: Int = 1, idBufferBatches: Int = maxConcurrentReqs * 4, + localEntries: List = emptyList(), onProgress: ((needSoFar: Int, downloaded: Int) -> Unit)? = null, onEvent: (Event) -> Unit, ): NegentropyOrFetchResult = @@ -356,6 +365,7 @@ suspend fun INostrClient.negentropySyncOrFetch( idleTimeoutMs = idleTimeoutMs, reconcileConcurrency = reconcileConcurrency, idBufferBatches = idBufferBatches, + localEntries = localEntries, onProgress = onProgress, onEvent = onEvent, ) @@ -391,6 +401,7 @@ private suspend fun INostrClient.syncPipeline( maxConcurrentReqs: Int, reconcileConcurrency: Int, idBufferBatches: Int, + localEntries: List, onWindow: () -> Unit, onNeed: (Int) -> Unit, deliver: suspend (Event) -> Unit, @@ -409,11 +420,15 @@ private suspend fun INostrClient.syncPipeline( } } + // reconcileWindows needs the local set sorted by createdAt (it binary-searches + // each window's slice). Empty/singleton sets are already trivially sorted. + val sortedLocal = if (localEntries.size > 1) localEntries.sortedBy { it.createdAt } else localEntries + reconcileWindows( clients = listOf(this@syncPipeline), relay = relay, filter = filter, - localEntries = emptyList(), + localEntries = sortedLocal, idleTimeoutMs = idleTimeoutMs, batchSize = fetchBatch, reconcileConcurrency = reconcileConcurrency, diff --git a/relayBench/plans/2026-07-04-sync-throughput-1m.md b/relayBench/plans/2026-07-04-sync-throughput-1m.md index bed5fc03a1..ef8c9638de 100644 --- a/relayBench/plans/2026-07-04-sync-throughput-1m.md +++ b/relayBench/plans/2026-07-04-sync-throughput-1m.md @@ -129,16 +129,35 @@ is the direct answer to "why aren't we using negentropy for geode too": we should, and over plain REQ geode structurally *cannot* finish a 1M pull from strfry. -### Recommendation (needs a design call) +### Implemented: MirrorWorker now mirrors strfry's two-phase model -A production geode that backfills a large foreign relay from empty should use -**NIP-77 negentropy** for the initial pull (reconcile → client-paced fetch), -not the live-tail REQ path. geode already has the negentropy client/server; -wiring a "reconcile then fetch the diff in bounded batches" backfill into -`MirrorWorker` would make strfry→geode both robust and directly comparable to -strfry→strfry. (A cheaper stop-gap — bound MirrorWorker's intake per-upstream -with a dedicated reader — fixes the OOM but not the strfry-kills-slow-clients -wall; only negentropy fixes both.) +`MirrorWorker` gained a NIP-77 **"sync" catch-up** phase (geode's equivalent of +`strfry sync --dir down`) that runs once per down/both upstream before the live +REQ tail: + +- **Catch-up** reconciles the local set against the upstream over the historical + `[now - backfill_seconds, now]` window and downloads only the diff, via the + ready-made `INostrClient.negentropySyncOrFetch` — client-paced, so strfry can't + overrun/kill us, and it **completes** the bulk pull. A small + `localEntries` param was added to the public `negentropySync`/ + `negentropySyncOrFetch` so the reconcile diffs against what we already hold + (incremental, like `strfry sync`) instead of re-downloading. +- **Either mode, transparently**: `negentropySyncOrFetch` auto-falls back to + paged REQ for an upstream that doesn't speak NIP-77 — no config toggle. +- **Live tail** unchanged: the REQ subscription now starts at `now` when catch-up + is on (history is the sync's job); the two windows overlap at `now` and the + store's unique-id constraint dedups the seam. + +Same vocabulary as strfry throughout — one `[[mirror]]` entry, one `dir` +(down/up/both) driving both phases; negentropy-vs-REQ is an internal transport +detail. The geode binary opts in (`Main` passes `negentropyBackfill = true` + +the store); the `MirrorWorker` default stays off so existing live-REQ tests are +unchanged. See `MirrorNegentropyCatchUpTest` (catch-up isolated from the live +tail by preloading *historical* events a live-only sub can't deliver). + +Remaining follow-ups: negentropy for the **up** direction (currently REQ replay); +`liveNegentropySnapshot`-based local enumeration for very large mirrors; and +optionally bounding the live-tail intake per-upstream. Separately worth a look: geode's real-content ingest *decays* from ~11k→~7k ev/s as the in-memory store grows to a few hundred k — expected B-tree/FTS