diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/SyncCommand.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/SyncCommand.kt index 439da7bc2a..9f878186d0 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/SyncCommand.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/SyncCommand.kt @@ -26,28 +26,27 @@ import com.vitorpamplona.amethyst.cli.DataDir import com.vitorpamplona.amethyst.cli.Output import com.vitorpamplona.quartz.nip01Core.core.Event import com.vitorpamplona.quartz.nip01Core.core.HexKey -import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper +import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.NegentropySyncException +import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.negentropyReconcile import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer -import com.vitorpamplona.quartz.nip01Core.relay.normalizer.toHttp -import com.vitorpamplona.quartz.nip77Negentropy.NegErrMessage -import com.vitorpamplona.quartz.nip77Negentropy.NegMsgMessage -import com.vitorpamplona.quartz.nip77Negentropy.NegentropySession +import com.vitorpamplona.quartz.nip01Core.store.IdAndTime import kotlinx.coroutines.channels.Channel -import kotlinx.coroutines.channels.Channel.Factory.UNLIMITED -import kotlinx.coroutines.withTimeoutOrNull -import okhttp3.OkHttpClient -import okhttp3.Request -import okhttp3.Response -import okhttp3.WebSocket -import okhttp3.WebSocketListener +import kotlinx.coroutines.coroutineScope +import kotlinx.coroutines.joinAll +import kotlinx.coroutines.launch +import java.util.concurrent.atomic.AtomicInteger /** * `amy sync --relay URL [filter flags] [--down] [--up] [--timeout SECS]` * * NIP-77 Negentropy set-reconciliation between the local event store and a - * relay (nak's `sync`, adapted to amy's local-store model). First it - * negotiates the symmetric difference under [filter], then closes the loop: + * relay (nak's `sync`, adapted to amy's local-store model). The protocol + * itself is quartz's [negentropyReconcile]: it pins the relay with a + * keep-alive subscription, splits the filter by `created_at` window whenever + * the relay caps the set (strfry `max_sync_events`), and streams the two + * directions of the diff as each round completes. This command closes the + * loop on that stream: * * --down (default) download events the relay has and we lack (REQ by id) * --up upload events we have and the relay lacks (EVENT) @@ -55,13 +54,30 @@ import okhttp3.WebSocketListener * Pass both for a full bidirectional sync. The filter flags are the same as * `fetch`/`subscribe`; an empty filter reconciles the whole store. * - * Thin assembly only: the negentropy protocol lives in quartz - * (`NegentropySession`); this file drives the WebSocket round-trips and - * reuses `Context.drain` / `Context.publish` for the NIP-01 follow-up. + * Both directions are pipelined with the reconcile: need-id batches feed + * [DOWNLOAD_WORKERS] concurrent by-id REQ drains and have-ids feed a single + * uploader, so downloads and uploads overlap the remaining reconcile rounds + * instead of waiting for the full diff. Every downloaded event funnels + * through `Context.drain`'s verify-and-store path, unchanged. + * + * Thin assembly only: the windowing, streaming, and back-pressure live in + * quartz (`negentropyReconcile`); this file only routes ids to + * `Context.drain` / `Context.publish`. */ object SyncCommand { private const val ID_CHUNK = 500 + /** + * Concurrent by-id download REQs. With [RECONCILE_CONCURRENCY] NEG + * sessions and the keep-alive, peak concurrent subscriptions on the + * relay are DOWNLOAD_WORKERS + RECONCILE_CONCURRENCY + 1 = 7 — well + * under the common NIP-11 `max_subscriptions` floor of 20. + */ + private const val DOWNLOAD_WORKERS = 4 + + /** Overlapped `created_at`-window reconciles after an over-cap split. */ + private const val RECONCILE_CONCURRENCY = 2 + suspend fun run( dataDir: DataDir, rest: Array, @@ -83,127 +99,79 @@ object SyncCommand { ctx.prepare() val localEvents = ctx.store.query(filter) val localById = localEvents.associateBy { it.id } + val localEntries = localEvents.map { IdAndTime(it.createdAt, it.id) } - val diff = - negotiate(relay.toHttp(), filter, localEvents, timeoutMs) - ?: return Output.error("timeout", "no negentropy response from ${relay.url} within ${timeoutMs}ms") - if (diff.error != null) { - return Output.error("sync_error", diff.error) - } + val downloaded = AtomicInteger(0) + val uploaded = AtomicInteger(0) - // needIds = relay has, we lack; haveIds = we have, relay lacks. - var downloaded = 0 - if (down && diff.needIds.isNotEmpty()) { - diff.needIds.chunked(ID_CHUNK).forEach { chunk -> - val got = ctx.drain(mapOf(relay to listOf(Filter(ids = chunk))), timeoutMs) - downloaded += got.size + val result = + try { + coroutineScope { + // needIds = relay has, we lack; haveIds = we have, relay lacks. + // Bounded so a slow download back-pressures the reconcile + // rounds instead of piling ids up in memory. + val needBatches = Channel>(DOWNLOAD_WORKERS * 2) + // Unbounded is fine here: have-ids reference events we already + // hold locally, so memory is bounded by the local set. + val haveBatches = Channel>(Channel.UNLIMITED) + + val downloaders = + List(DOWNLOAD_WORKERS) { + launch { + for (batch in needBatches) { + val got = ctx.drain(mapOf(relay to listOf(Filter(ids = batch))), timeoutMs) + downloaded.addAndGet(got.size) + } + } + } + val uploader = + launch { + for (batch in haveBatches) { + for (id in batch) { + val ev = localById[id] ?: continue + val ack = ctx.publish(ev, setOf(relay)) + if (ack.values.any { it }) uploaded.incrementAndGet() + } + } + } + + val reconcile = + try { + ctx.client.negentropyReconcile( + relay = relay, + filter = filter, + localEntries = localEntries, + batchSize = ID_CHUNK, + idleTimeoutMs = timeoutMs, + reconcileConcurrency = RECONCILE_CONCURRENCY, + onHaveIds = if (up) { batch -> haveBatches.send(batch) } else null, + onNeedIds = { batch -> if (down) needBatches.send(batch) }, + ) + } finally { + needBatches.close() + haveBatches.close() + } + + downloaders.joinAll() + uploader.join() + reconcile + } + } catch (e: NegentropySyncException) { + return Output.error("sync_error", e.message ?: "negentropy sync failed") } - } - - var uploaded = 0 - if (up && diff.haveIds.isNotEmpty()) { - diff.haveIds.forEach { id -> - val ev = localById[id] ?: return@forEach - val ack = ctx.publish(ev, setOf(relay)) - if (ack.values.any { it }) uploaded++ - } - } Output.emit( mapOf( "relay" to relay.url, "local_events" to localEvents.size, - "rounds" to diff.rounds, - "need" to diff.needIds.size, - "have" to diff.haveIds.size, - "downloaded" to downloaded, - "uploaded" to uploaded, + "windows" to result.windows, + "need" to result.needCount, + "have" to result.haveCount, + "downloaded" to downloaded.get(), + "uploaded" to uploaded.get(), ), ) return 0 } } - - private data class Diff( - val haveIds: List, - val needIds: List, - val rounds: Int, - val error: String?, - ) - - /** - * Drive one NIP-77 reconciliation over a raw WebSocket and return the - * symmetric difference. Mirrors geode's interop sync driver: the protocol - * state machine is [NegentropySession]; this only shuttles frames. - */ - private suspend fun negotiate( - httpUrl: String, - filter: Filter, - localEvents: List, - timeoutMs: Long, - subId: String = "amy-sync", - maxRounds: Int = 64, - ): Diff? { - val incoming = Channel(UNLIMITED) - val client = OkHttpClient.Builder().build() - val ws = - client.newWebSocket( - Request.Builder().url(httpUrl).build(), - object : WebSocketListener() { - override fun onMessage( - webSocket: WebSocket, - text: String, - ) { - incoming.trySend(text) - } - - override fun onClosing( - webSocket: WebSocket, - code: Int, - reason: String, - ) { - incoming.close() - } - - override fun onFailure( - webSocket: WebSocket, - t: Throwable, - response: Response?, - ) { - incoming.close(t) - } - }, - ) - - return try { - withTimeoutOrNull(timeoutMs) { - val session = NegentropySession.fromEvents(subId, filter, localEvents, frameSizeLimit = 0) - ws.send(OptimizedJsonMapper.toJson(session.open())) - - val have = mutableSetOf() - val need = mutableSetOf() - var rounds = 0 - while (rounds < maxRounds) { - val raw = incoming.receive() - rounds++ - when (val msg = OptimizedJsonMapper.fromJsonToMessage(raw)) { - is NegErrMessage -> return@withTimeoutOrNull Diff(have.toList(), need.toList(), rounds, "${msg.subId}: ${msg.reason}") - is NegMsgMessage -> { - val r = session.processMessage(msg.message) - have += r.haveIds - need += r.needIds - if (r.isComplete()) { - return@withTimeoutOrNull Diff(have.toList(), need.toList(), rounds, null) - } - ws.send(OptimizedJsonMapper.toJson(r.nextCmd!!)) - } - else -> {} // ignore NOTICE/etc. and keep waiting - } - } - Diff(have.toList(), need.toList(), rounds, "did not converge in $maxRounds rounds") - } - } finally { - ws.close(1000, "amy-sync-done") - } - } } diff --git a/quartz/plans/2026-07-02-nostrclient-receiver-perf.md b/quartz/plans/2026-07-02-nostrclient-receiver-perf.md index 97062c0234..68e423b3b6 100644 --- a/quartz/plans/2026-07-02-nostrclient-receiver-perf.md +++ b/quartz/plans/2026-07-02-nostrclient-receiver-perf.md @@ -440,6 +440,18 @@ now available standalone in `NostrClientNegentropySyncExt`: ways, identical sets, batch streaming, since/until window slicing) — 49 negentropy tests green. +**`amy sync` migrated onto it (2026-07-03):** the CLI's hand-rolled raw +WebSocket negotiate loop (single un-windowed session; a strfry +`max_sync_events` overflow was a hard error) is gone. `SyncCommand` now +calls `negentropyReconcile` on amy's own `NostrClient` and pipelines the +loop-closing: need-id batches feed 4 concurrent by-id `drain`s and have-ids +feed an uploader, both overlapping the remaining reconcile rounds (peak 7 +subs on the relay, under the common cap of 20). Downloads still funnel +through amy's verify-and-store path. Output field `rounds` (protocol +round-trips) became `windows` (created_at splits). Verified end-to-end +against two embedded geode relays: down-only 25/25, up-only 5/5, and +bidirectional re-runs converge to a zero diff. + ### Answer for "10M events from one relay, fastest" At nosfabrica's measured page cadence, one connection ≈ 3.7k events/s → 10M @@ -515,6 +527,13 @@ lock-free wins 3.0–4.5×** (e.g. 127ms → 28ms for 128k frames). The lock was serializing the concurrent hit path exactly the way the old global PoolRequests lock did; ConcurrentHashMap removes it. +**Adopted in all four front-end clients (2026-07-03):** the Android app's +pool (`AppModules.kt`), the Android crawl client +(`AccountViewModel.buildCrawlClient` — Event Sync / Cashu discovery, the +duplicate-heaviest path), the desktop `RelayConnectionManager`, and amy's +`Context`. Each passes `CachingEventDecoder()` at `NostrClient` +construction; no other behavior change. + ### Bounded per-connection receive buffer (REVERTED — deliberate design choice) A 4096-frame bound on the reader→consumer channel was tried (backpressure