mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 03:38:23 +00:00
feat(geode): mirror strfry's two-phase model — NIP-77 sync catch-up + live REQ tail
geode's MirrorWorker mirrored `strfry router` (live REQ streaming) but had no `strfry sync` equivalent, so backfilling a large foreign relay from empty could not complete: a plain REQ dump of the history overruns the sink and strfry kills the slow client at its maxPendingOutboundBytes cap (see relayBench/plans/2026-07-04-sync-throughput-1m.md). MirrorWorker now runs a one-shot NIP-77 "sync" catch-up per down/both upstream before the live tail, using strfry's own vocabulary — one `[[mirror]]` entry, one `dir` driving both phases: - Catch-up reconciles the local set against the upstream over the [now - backfill_seconds, now] window and downloads only the diff via the existing INostrClient.negentropySyncOrFetch — client-paced (strfry can't overrun us) and it completes the pull. Reconcile-against-local means a warm restart re-fetches nothing it already holds, like `strfry sync`. - Either mode, transparently: negentropySyncOrFetch auto-falls back to paged REQ for an upstream without NIP-77 — no config toggle. - Live REQ tail unchanged; it starts at `now` when catch-up is on (history is the sync's job). The windows overlap at `now`; the store's unique-id constraint dedups the seam. Changes: - quartz: add a backward-compatible `localEntries` param to the public negentropySync / negentropySyncOrFetch (default empty = prior behavior) so the reconcile diffs against a caller-supplied local set. - geode MirrorWorker: `runCatchUp()` (bounded, backpressured ingest; same trusted-scope re-check as the live path; failure is non-fatal). New `store` + `negentropyBackfill` ctor params; default off so existing live-REQ tests are unchanged. Main opts production in. - Test: MirrorNegentropyCatchUpTest isolates catch-up from the live tail by preloading historical events a live-only sub cannot deliver, then proves the post-boot event still arrives (3000 catch-up + 1 live = 3001). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_012EZeWww5TJnzBZKPoc6mvU
This commit is contained in:
@@ -260,7 +260,17 @@ fun main(args: Array<String>) {
|
||||
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`).
|
||||
|
||||
@@ -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<MirrorUpstream>,
|
||||
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<Event>(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
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 ─")
|
||||
}
|
||||
}
|
||||
+24
-9
@@ -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<IdAndTime> = 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<IdAndTime> = 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<IdAndTime> = 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<IdAndTime> = 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<IdAndTime>,
|
||||
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,
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user