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 4e3e6c842d..0a548e5b64 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.negentropyReconcileIds 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 @@ -329,24 +330,30 @@ 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 + // strfry's two-phase model, both directions (`strfry sync --dir + // both` + the live router): a one-shot NIP-77 "sync" closes the + // historical [initialSince, now] gap — down pulls what the upstream + // has and we lack, up pushes what we have and it lacks — then the + // live REQ subscription/session tails everything new. When a + // direction's catch-up is on, that direction's live window starts at + // `now` (history is the sync's job); the windows overlap at `now` and + // the store's/upstream's unique-id dedup absorbs the seam. + val catchUpDown = negentropyBackfill && up.direction != MirrorDirection.UP && up.backfillSeconds > 0 + val catchUpUp = negentropyBackfill && up.direction != MirrorDirection.DOWN && up.backfillSeconds > 0 + val downLiveSince = if (catchUpDown) since else initialSince + val upLiveSince = if (catchUpUp) since else initialSince if (up.direction != MirrorDirection.UP) { downSubs += startDown(i, up, scopedBase, downLiveSince, exchanged) } if (up.direction != MirrorDirection.DOWN) { - startUp(up, scopedBase.copy(since = initialSince), exchanged) + startUp(up, scopedBase.copy(since = upLiveSince), exchanged) } - if (catchUp) { - scope.launch { runCatchUp(up, scopedBase, initialSince, since) } + if (catchUpDown) { + scope.launch { runCatchUpDown(up, scopedBase, initialSince, since) } + } + if (catchUpUp) { + scope.launch { runCatchUpUp(up, scopedBase, initialSince, since, exchanged) } } } client.connect() @@ -399,7 +406,7 @@ class MirrorWorker( * 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( + private suspend fun runCatchUpDown( up: MirrorUpstream, scopedBase: Filter, initialSince: Long, @@ -458,13 +465,88 @@ class MirrorWorker( } catch (e: CancellationException) { throw e } catch (e: Throwable) { - Log.w("MirrorWorker") { "catch-up from ${up.url.url} failed (live tail continues): ${e.message}" } + Log.w("MirrorWorker") { "down catch-up from ${up.url.url} failed (live tail continues): ${e.message}" } } finally { handoff.close() consumer.join() } } + /** + * The up half of the NIP-77 "sync" — geode's `strfry sync --dir up`. + * Reconciles the local set against the upstream and PUSHES the events we + * hold that the upstream lacks (the reconcile's `have` ids) — the mirror of + * [runCatchUpDown]. Negentropy-only: if the upstream doesn't speak NIP-77 the + * reconcile throws and we log; the live up-session (started at `now`) keeps + * pushing new local events, so this is a non-fatal historical catch-up. + * + * Needs [store] to enumerate what we hold; without it there is nothing to + * reconcile against and the phase is skipped. + */ + private suspend fun runCatchUpUp( + up: MirrorUpstream, + scopedBase: Filter, + initialSince: Long, + until: Long, + exchanged: RecentIds?, + ) { + val localStore = store ?: return + val catchUpFilter = scopedBase.copy(since = initialSince, until = until) + // Our local set for this window is fixed; each round reconciles it + // against the upstream, which grows as we push, so the `have` diff + // shrinks to zero. + val localEntries = localStore.snapshotIdsForNegentropy(listOf(catchUpFilter)) + try { + var round = 0 + while (round < MAX_UP_SYNC_ROUNDS) { + // Reconcile (need ids are the down catch-up's job — ignore them); + // `have` = events we hold that the upstream still lacks. + val diff = + client.negentropyReconcileIds( + relay = up.url, + filter = catchUpFilter, + localEntries = localEntries, + ) + if (diff.haveIds.isEmpty()) { + Log.i("MirrorWorker") { "up catch-up to ${up.url.url}: converged after $round round(s)" } + return + } + // Publish the remaining local-only events. `client.publish`'s + // outbox is best-effort under a bulk burst (each publish also + // churns a reconnect), so instead of trusting one pass we + // re-reconcile next round and re-push only what didn't land — + // the reconcile is the delivery check, so the push converges + // to lossless. + var pushed = 0 + for (batch in diff.haveIds.chunked(HAVE_FETCH_BATCH)) { + for (event in localStore.query(Filter(ids = batch))) { + // Scope containment: a scoped upstream only receives + // in-scope events. Echo suppression: record the id so a + // BOTH mirror doesn't re-ingest its own push on the down + // sub (but always re-publish — a straggler stays in + // `exchanged` yet still needs delivering). + if (up.filter != null && !up.filter.match(event)) continue + exchanged?.add(event.id) + client.publish(event, setOf(up.url)) + pushed++ + } + delay(UP_PUBLISH_PACING_MS) + } + sentUp.addAndGet(pushed.toLong()) + Log.i("MirrorWorker") { "up catch-up to ${up.url.url}: round $round pushed $pushed (had ${diff.haveIds.size} to go)" } + round++ + // Let the upstream ingest + OK before the next reconcile, so the + // diff reflects what actually landed rather than what's in flight. + delay(UP_SYNC_SETTLE_MS) + } + Log.w("MirrorWorker") { "up catch-up to ${up.url.url}: did not fully converge in $MAX_UP_SYNC_ROUNDS rounds (live push continues)" } + } catch (e: CancellationException) { + throw e + } catch (e: Throwable) { + Log.w("MirrorWorker") { "up catch-up to ${up.url.url} failed (live push continues): ${e.message}" } + } + } + /** Down direction: subscribe to the upstream, ingest what it sends. */ private fun startDown( index: Int, @@ -616,5 +698,17 @@ class MirrorWorker( * behind `server.ingest` is the real limiter. */ const val CATCHUP_HANDOFF = 4_096 + + /** Ids per store read when loading up-catch-up events to publish. */ + const val HAVE_FETCH_BATCH = 500 + + /** Pause between up-catch-up publish batches so the outbox drains. */ + const val UP_PUBLISH_PACING_MS = 40L + + /** Max reconcile→push rounds before the up catch-up gives up converging. */ + const val MAX_UP_SYNC_ROUNDS = 8 + + /** Settle time between an up-push and the next verifying reconcile. */ + const val UP_SYNC_SETTLE_MS = 1_500L } } diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorNegentropyCatchUpTest.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorNegentropyCatchUpTest.kt index 1664e863a1..9bd32d24f5 100644 --- a/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorNegentropyCatchUpTest.kt +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/MirrorNegentropyCatchUpTest.kt @@ -93,6 +93,7 @@ class MirrorNegentropyCatchUpTest { } private suspend fun awaitCount( + store: EventStore, target: Int, timeoutMs: Long = 120_000, ): Int { @@ -101,7 +102,7 @@ class MirrorNegentropyCatchUpTest { var stable = 0 withTimeoutOrNull(timeoutMs) { while (true) { - val c = downstreamStore.count(Filter()) + val c = store.count(Filter()) reached = c if (c >= target) break if (c == last) { @@ -159,7 +160,7 @@ class MirrorNegentropyCatchUpTest { ).also { it.start() } // Phase 1: the catch-up must deliver every historical event. - val afterCatchUp = awaitCount(count) + val afterCatchUp = awaitCount(downstreamStore, 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 @@ -177,13 +178,68 @@ class MirrorNegentropyCatchUpTest { ) upstream.server.ingest(live, skipVerify = true) { } - val afterLive = awaitCount(count + 1) + val afterLive = awaitCount(downstreamStore, 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 ─") + println("─ MirrorNegentropyCatchUpTest(down): catch-up $count + live 1 = $afterLive delivered ─") + } + + /** + * The mirror image — `dir = up`, geode's `strfry sync --dir up`. The local + * geode holds a set of historical events; the upstream sink is empty. The + * up catch-up must reconcile and PUSH every local event to the sink. Uses a + * default (no-verify) sink so the synthetic-signature events are accepted. + */ + @Test + fun negentropyCatchUpPushesUp() = + runBlocking { + val count = System.getProperty("catchUpUpN")?.toInt() ?: 3_000 + val now = TimeUtils.now() + val sig = "f".repeat(128) + + // The LOCAL geode (downstream) holds the events; the sink (upstream) + // starts empty and must receive them all via the up push. + val history = + (0 until count).map { i -> + Event( + id = hex64(11, i), + pubKey = hex64(3, i % 500), + createdAt = now - 3_600 - (i % 1_000), + kind = 1, + tags = emptyArray(), + content = "u$i", + sig = sig, + ) + } + history.chunked(10_000).forEach { downstreamStore.batchInsert(it) } + assertEquals(count, downstreamStore.count(Filter()), "local geode preloaded") + assertEquals(0, upstreamStore.count(Filter()), "sink starts empty") + + // The sink is the relay geode dials; it must accept what we push. + 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 = false, + backfillSeconds = 86_400, + direction = MirrorDirection.UP, + ), + ), + server = downstream.server, + store = downstreamStore, + negentropyBackfill = true, + ).also { it.start() } + + val pushed = awaitCount(upstreamStore, count) + assertEquals(count, pushed, "up catch-up failed to push ${count - pushed} of $count events") + println("─ MirrorNegentropyCatchUpTest(up): pushed $pushed/$count to the sink via negentropy ─") } } diff --git a/relayBench/plans/2026-07-04-sync-throughput-1m.md b/relayBench/plans/2026-07-04-sync-throughput-1m.md index ef8c9638de..d41c52af99 100644 --- a/relayBench/plans/2026-07-04-sync-throughput-1m.md +++ b/relayBench/plans/2026-07-04-sync-throughput-1m.md @@ -148,16 +148,37 @@ REQ tail: is on (history is the sync's job); the two windows overlap at `now` and the store's unique-id constraint dedups the seam. +**Both directions**, matching `strfry sync --dir both` (which strfry's source +confirms is bidirectional negentropy — `doUp = both||up`, `doDown = both||down`): + +- **down** catch-up pulls what the upstream has and we lack + (`negentropySyncOrFetch`, with paged fallback), then the live REQ sub tails. +- **up** catch-up reconciles and pushes what we have and the upstream lacks + (`negentropyReconcileIds`'s `have` ids → `client.publish`), then the live + up-session tails. Negentropy-only (no paged fallback needed — the live + up-session covers a non-NIP-77 upstream). The push runs as a **reconcile→push + convergence loop**: `client.publish`'s outbox is best-effort under a bulk + burst (each publish also churns a reconnect; measured ~1–2% dropped per pass), + so each round re-reconciles — the reconcile *is* the delivery check — and + re-pushes only the stragglers until the `have` diff is empty. Observed on the + 3000-event up test: 3000 → 69 → 2 → 0 across 3 rounds, lossless. + 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). +unchanged. See `MirrorNegentropyCatchUpTest`: the down test isolates catch-up +from the live tail by preloading *historical* events a live-only sub can't +deliver (3000 + 1 live); the up test pushes 3000 local events to an empty sink. -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. +Structural note: strfry keeps these as *separate commands* (`sync` = negentropy +both-ways; `router`/`stream` = REQ live both-ways). geode folds both into one +`MirrorWorker` under a single `dir` — more integrated, same semantics. + +Remaining follow-ups: `liveNegentropySnapshot`-based local enumeration for very +large mirrors (the up path's `snapshotIdsForNegentropy` currently scans); a +single-pass reconcile for `dir=both` (today it runs one reconcile per +direction); 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