From 677c0ee2074f274af5a5ad40b8a183e6206b66c5 Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 8 Jul 2026 14:25:37 +0000 Subject: [PATCH] refactor: deletion sync as a post-settle residual pass (both directions, O(residual)) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Replace the per-need-event fetch (which pulled the whole need set just to read metadata — an O(db) regression on large syncs) with a second reconcile pass over the residual, per the "settle, then diff, then explain what didn't converge" idea. Pass 1 is the plain content sync again (drain needs, publish haves) — zero deletion overhead. Pass 2+ re-reconciles; the leftover diff is exactly the deletion mismatches, and only that (tiny) set is fetched: - residual need (relay has it, we still lack it after --down) = we deleted it → publish our covering deletion up so the relay drops it; - residual have (we have it, relay still lacks it after --up) = the relay deleted it → pull the relay's covering kind-5 down and apply locally (vanish is NOT auto-applied on pull — account-wide blast radius). Loops until a round resolves nothing (converges + self-verifies). So `amy sync` makes the relay honor our deletions; `--up` makes us honor the relay's; `--up --down` converges both ways. Cost is one cheap reconcile + the residual regardless of database size — the large-DB bottleneck is gone by construction, not by heuristics. quartz: deletionsCovering is now source-agnostic (takes a query lambda) so the same coverage rule runs against the local store (up) or the relay (down); the IEventStore overload is the local convenience. Tests: DeletionSyncTest gains the down-direction end-to-end (relay deleted → local removes) alongside the up-direction and the per-form unit cases. Co-Authored-By: Claude Opus 4.8 Claude-Session: https://claude.ai/code/session_01JgL1WTV4Hkp2uuXcUHCHGt --- .../amethyst/cli/commands/SyncCommand.kt | 169 ++++++++++++------ .../vitorpamplona/geode/DeletionSyncTest.kt | 42 +++++ .../nip01Core/store/EventStoreDeletionsExt.kt | 38 ++-- 3 files changed, 180 insertions(+), 69 deletions(-) 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 333c8a2677..540b5c9894 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 @@ -29,15 +29,16 @@ import com.vitorpamplona.quartz.nip01Core.core.HexKey import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.NegentropySyncException import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAll import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.negentropyReconcile +import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.negentropyReconcileIds import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer import com.vitorpamplona.quartz.nip01Core.store.IdAndTime import com.vitorpamplona.quartz.nip01Core.store.deletionsCovering +import com.vitorpamplona.quartz.nip09Deletions.DeletionEvent import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.coroutineScope import kotlinx.coroutines.joinAll import kotlinx.coroutines.launch -import java.util.concurrent.ConcurrentHashMap import java.util.concurrent.atomic.AtomicInteger /** @@ -57,25 +58,29 @@ import java.util.concurrent.atomic.AtomicInteger * Pass both for a full bidirectional sync. The filter flags are the same as * `fetch`/`subscribe`; an empty filter reconciles the whole store. * - * Deletion propagation is deliberately narrow (on by default; disable with - * `--no-sync-deletions`): for the events the relay HAS that we LACK — the reconcile's - * need set — we publish up the local deletions that would make the relay remove them, - * and only those. That covers a NIP-09 kind-5 targeting the event by id (`e` tag) or - * by address (`a` tag, cutoff-checked), and a NIP-62 kind-62 vanish for the event's - * author that targets this relay. The need events are fetched only for their metadata - * (author/address/created_at); nothing is pulled down or applied locally, so it can - * never over-delete this store, and the need set already bounds it (no author scoping). - * See [com.vitorpamplona.quartz.nip01Core.store.deletionsCovering]. + * Deletion propagation (on by default; disable with `--no-sync-deletions`) is a + * **second pass over the residual**, not per-event work in the content pass — so it + * costs the same whether the database is tiny or huge. After the content settle, a + * re-reconcile's leftover diff is (barring races) exactly the events a deletion kept + * from converging: * - * 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. + * - a residual **need** (relay has it, we still lack it after `--down` tried to + * download) = we deleted it → publish OUR covering deletion up so the relay drops it; + * - a residual **have** (we have it, relay still lacks it after `--up` tried to upload) + * = the relay deleted it → pull the relay's covering kind-5 down and apply it locally. * - * Thin assembly only: the windowing, streaming, and back-pressure live in - * quartz (`negentropyReconcile`); this file only routes ids to - * `Context.drain` / `Context.publish`. + * Coverage is any way a deletion reaches an event ([deletionsCovering]): a NIP-09 kind-5 + * by id (`e`) or address (`a`, cutoff-checked), or a NIP-62 vanish targeting this relay + * (up direction only — a pulled vanish is not auto-applied, its blast radius being the + * whole account). The residual is small (only real deletion mismatches), so only it is + * fetched — never the whole need set. The loop repeats until a round resolves nothing. + * So `amy sync` (default `--down`) makes the relay honor your deletions; `--up` makes + * your store honor the relay's; `--up --down` converges both ways. + * + * Content is pipelined with the reconcile: need-id batches feed [DOWNLOAD_WORKERS] + * concurrent by-id REQ drains and have-ids feed a single uploader. 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 @@ -91,6 +96,15 @@ object SyncCommand { /** Overlapped `created_at`-window reconciles after an over-cap split. */ private const val RECONCILE_CONCURRENCY = 2 + /** + * Cap on deletion-settle rounds. Each round resolves the residual it can and + * re-reconciles; a healthy sync converges in 1–2 (round N sends/applies, round + * N+1 confirms empty). The cap only bounds pathological non-convergence (e.g. a + * relay that refuses a deletion), which the "resolved nothing → stop" check + * normally catches first. + */ + private const val MAX_DELETION_ROUNDS = 4 + suspend fun run( dataDir: DataDir, rest: Array, @@ -106,15 +120,6 @@ object SyncCommand { // Default direction is download; --up adds upload. val up = args.bool("up") val down = args.bool("down") || !up - // Deletion propagation (on by default; --no-sync-deletions disables). Scope is - // exactly: for the events the relay HAS that we LACK (the reconcile's need set), - // publish the local deletions that would make the relay remove them — an id- or - // address-based kind-5, or a kind-62 vanish that targets this relay. Only those - // deletions, nothing else (not other deletions by the same author). We fetch the - // need events (not to keep — `fetchAll` neither verifies nor stores) only - // to learn their author/address/created_at so [deletionsCovering] can tell which - // of our deletions actually apply. Nothing is pulled down or applied locally, so - // this can never over-delete the local store. val syncDeletions = !args.bool("no-sync-deletions") val filter = RawEventSupport.buildFilter(args) @@ -126,42 +131,23 @@ object SyncCommand { val downloaded = AtomicInteger(0) val uploaded = AtomicInteger(0) - val deletionsSent = AtomicInteger(0) - // Deduplicate published deletions across the concurrent need workers: one - // deletion often covers several need events. - val sentDeletions = ConcurrentHashMap.newKeySet() + // ── Pass 1: content settle — download needs, upload haves. No deletion + // logic, so a plain sync costs exactly what it always did. val result = try { coroutineScope { // needIds = relay has, we lack; haveIds = we have, relay lacks. - // Bounded so a slow worker 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 needWorkers = + val downloaders = List(DOWNLOAD_WORKERS) { launch { for (batch in needBatches) { - // Fetch the need events once (no verify/store — we only - // need their metadata to decide which deletions apply). - val events = ctx.client.fetchAll(relay, Filter(ids = batch), timeoutMs) - // Push up the deletions that would remove them from the relay. - if (syncDeletions) { - for (del in ctx.store.deletionsCovering(events, relay)) { - if (sentDeletions.add(del.id) && ctx.publish(del, setOf(relay)).values.any { it }) { - deletionsSent.incrementAndGet() - } - } - } - // Download the rest into the local store; anything we - // deleted is rejected by the store's own tombstone. - if (down) { - for (event in events) if (ctx.verifyAndStore(event)) downloaded.incrementAndGet() - } + // drain verifies + stores; anything we deleted is + // rejected by our own tombstone and stays a "need". + downloaded.addAndGet(ctx.drain(mapOf(relay to listOf(Filter(ids = batch))), timeoutMs).size) } } } @@ -185,16 +171,14 @@ object SyncCommand { idleTimeoutMs = timeoutMs, reconcileConcurrency = RECONCILE_CONCURRENCY, onHaveIds = if (up) { batch -> haveBatches.send(batch) } else null, - // Fetch need events when we either download them or need - // their metadata to decide which deletions to send. - onNeedIds = { batch -> if (down || syncDeletions) needBatches.send(batch) }, + onNeedIds = { batch -> if (down) needBatches.send(batch) }, ) } finally { needBatches.close() haveBatches.close() } - needWorkers.joinAll() + downloaders.joinAll() uploader.join() reconcile } @@ -202,6 +186,75 @@ object SyncCommand { return Output.error("sync_error", e.message ?: "negentropy sync failed") } + // ── Pass 2+: deletion settle. After the content pass, a re-reconcile's + // residual is (barring races) exactly the events a deletion kept from moving: + // - a residual NEED (relay has it, we still lack it after trying to download) + // = we deleted it → publish OUR covering deletion up so the relay drops it; + // - a residual HAVE (we have it, relay still lacks it after trying to upload) + // = the relay deleted it → pull the relay's covering kind-5 down and apply. + // The residual is tiny (only real deletion mismatches), so this is cheap no + // matter how large the database is — we only fetch metadata for the residual, + // never the whole need set. Loop until a round resolves nothing (converged) or + // we hit the round cap. Best-effort: a failed reconcile here never fails the + // command — the content sync already succeeded. + var deletionsUp = 0 + var deletionsDown = 0 + var deletionRounds = 0 + if (syncDeletions && (down || up)) { + val sentUp = HashSet() + val appliedDown = HashSet() + try { + while (deletionRounds < MAX_DELETION_ROUNDS) { + deletionRounds++ + val diff = + ctx.client.negentropyReconcileIds( + relay = relay, + filter = filter, + localEntries = ctx.store.snapshotIdsForNegentropy(listOf(filter)), + batchSize = ID_CHUNK, + idleTimeoutMs = timeoutMs, + reconcileConcurrency = RECONCILE_CONCURRENCY, + ) + var resolved = 0 + + // residual needs → send our deletions up (bounded: --down settled + // every need we don't have a deletion for). + if (down) { + for (chunk in diff.needIds.chunked(ID_CHUNK)) { + val events = ctx.client.fetchAll(relay, Filter(ids = chunk), timeoutMs) + for (del in ctx.store.deletionsCovering(events, relay)) { + if (sentUp.add(del.id) && ctx.publish(del, setOf(relay)).values.any { it }) { + deletionsUp++ + resolved++ + } + } + } + } + + // residual haves → apply the relay's deletions locally (bounded: + // --up settled every have the relay didn't delete). Only precise + // kind-5 deletions are pulled down; a kind-62 vanish is NOT + // auto-applied (its blast radius is the whole account). + if (up) { + for (chunk in diff.haveIds.chunked(ID_CHUNK)) { + val ourEvents = ctx.store.query(Filter(ids = chunk)) + val relayDeletions = deletionsCovering(ourEvents, relay) { f -> ctx.client.fetchAll(relay, f, timeoutMs) } + for (del in relayDeletions.filterIsInstance()) { + if (appliedDown.add(del.id) && ctx.verifyAndStore(del)) { + deletionsDown++ + resolved++ + } + } + } + } + + if (resolved == 0) break + } + } catch (e: NegentropySyncException) { + // content already synced; deletion convergence is best-effort. + } + } + Output.emit( mapOf( "relay" to relay.url, @@ -211,7 +264,9 @@ object SyncCommand { "have" to result.haveCount, "downloaded" to downloaded.get(), "uploaded" to uploaded.get(), - "deletions_sent" to deletionsSent.get(), + "deletions_sent_up" to deletionsUp, + "deletions_applied_down" to deletionsDown, + "deletion_rounds" to deletionRounds, ), ) return 0 diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/DeletionSyncTest.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/DeletionSyncTest.kt index bd392fcfff..90327e529a 100644 --- a/geode/src/test/kotlin/com/vitorpamplona/geode/DeletionSyncTest.kt +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/DeletionSyncTest.kt @@ -123,6 +123,7 @@ class DeletionSyncTest : RelayClientTest() { // ---- end-to-end through the relay ---------------------------------------- + // UP direction: we deleted it, the relay still has it → send our deletion up. @Test fun sendsCoveringDeletionSoRelayRemovesTheNote() = runBlocking { @@ -154,4 +155,45 @@ class DeletionSyncTest : RelayClientTest() { "relay applied the pushed deletion and removed the note", ) } + + // DOWN direction: the relay deleted it, we still have it → pull the relay's deletion + // down and apply it locally (the residual-have resolution). + @Test + fun appliesRelaysDeletionSoLocalRemovesTheNote() = + runBlocking { + val target = note("delete me down") + val deletion = signer.sign(DeletionEvent.build(listOf(target), createdAt = target.createdAt + 1)) + + // Relay already applied the deletion → holds only the kind-5. + defaultRelay.preload(listOf(target, deletion)) + assertTrue(defaultRelay.store.query(Filter(ids = listOf(target.id))).isEmpty(), "relay deleted the note") + + // Local still holds the note (never saw the deletion). + val local = hub.getOrCreate(RelayUrlNormalizer.normalize("ws://local-down/")) + local.preload(listOf(target)) + assertEquals(1, local.store.query(Filter(ids = listOf(target.id))).size) + + // Reconcile → the note is a HAVE (we have it, the relay lacks it). + val diff = + withTimeout(20_000) { + client.negentropyReconcileIds( + relay = defaultRelayUrl, + filter = Filter(kinds = listOf(1)), + localEntries = listOf(IdAndTime(target.createdAt, target.id)), + ) + } + assertEquals(setOf(target.id), diff.haveIds.toSet()) + + // What SyncCommand does for the down direction: take our have events, ask the + // RELAY which of ITS deletions cover them, and apply those locally. + val ourEvents = local.store.query(Filter(ids = diff.haveIds)) + val relayDeletions = deletionsCovering(ourEvents, defaultRelayUrl) { f -> defaultRelay.store.query(f) } + assertEquals(listOf(deletion.id), relayDeletions.map { it.id }) + relayDeletions.filterIsInstance().forEach { local.store.insert(it) } + + assertTrue( + local.store.query(Filter(ids = listOf(target.id))).isEmpty(), + "local applied the pulled deletion and removed the note", + ) + } } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/EventStoreDeletionsExt.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/EventStoreDeletionsExt.kt index 2e07ee30bf..0c5a8203dc 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/EventStoreDeletionsExt.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/EventStoreDeletionsExt.kt @@ -50,23 +50,31 @@ private fun addressValue(event: Event): String { * `relay` tags name the URL or `ALL_RELAYS`), issued after that event (a vanish * deletes `created_at < vanish.created_at`). * - * Deduped by event id; a single deletion covering several server events is returned once. + * Deduped by event id; a single deletion covering several events is returned once. + * + * [query] is where the deletions are looked up — it is source-agnostic on purpose, so + * the same coverage rule runs in both sync directions: + * - **up** (send our deletions): `events` are the relay's, `query` is the local store — + * which of OUR deletions would delete what the relay still holds. + * - **down** (apply the relay's deletions): `events` are ours, `query` fetches from the + * relay — which of the RELAY'S deletions would delete what we still hold. */ -suspend fun IEventStore.deletionsCovering( - serverEvents: List, +suspend fun deletionsCovering( + events: List, relay: NormalizedRelayUrl, + query: suspend (Filter) -> List, ): List { - if (serverEvents.isEmpty()) return emptyList() + if (events.isEmpty()) return emptyList() val covering = LinkedHashMap() - // 1. id-based NIP-09: a kind-5 `e`-tagging a server id. - query(Filter(kinds = listOf(DeletionEvent.KIND), tags = mapOf("e" to serverEvents.map { it.id }))) + // 1. id-based NIP-09: a kind-5 `e`-tagging an event's id. + query(Filter(kinds = listOf(DeletionEvent.KIND), tags = mapOf("e" to events.map { it.id }))) .forEach { covering[it.id] = it } - // 2. address-based NIP-09: a kind-5 `a`-tagging a server event's coordinate, cutoff-checked. - val byAddress = serverEvents.filter { it.kind.isAddressable() || it.kind.isReplaceable() }.groupBy(::addressValue) + // 2. address-based NIP-09: a kind-5 `a`-tagging an event's coordinate, cutoff-checked. + val byAddress = events.filter { it.kind.isAddressable() || it.kind.isReplaceable() }.groupBy(::addressValue) if (byAddress.isNotEmpty()) { - query(Filter(kinds = listOf(DeletionEvent.KIND), tags = mapOf("a" to byAddress.keys.toList()))) + query(Filter(kinds = listOf(DeletionEvent.KIND), tags = mapOf("a" to byAddress.keys.toList()))) .forEach { del -> if (del !is DeletionEvent) return@forEach for (addr in del.deleteAddresses()) { @@ -79,12 +87,18 @@ suspend fun IEventStore.deletionsCovering( } } - // 3. NIP-62 vanish: a kind-62 by a server author, targeting this relay, issued after the event. - query(Filter(kinds = listOf(RequestToVanishEvent.KIND), authors = serverEvents.mapTo(HashSet()) { it.pubKey }.toList())) + // 3. NIP-62 vanish: a kind-62 by an event's author, targeting this relay, issued after it. + query(Filter(kinds = listOf(RequestToVanishEvent.KIND), authors = events.mapTo(HashSet()) { it.pubKey }.toList())) .forEach { vanish -> if (vanish !is RequestToVanishEvent || !vanish.shouldVanishFrom(relay)) return@forEach - if (serverEvents.any { it.pubKey == vanish.pubKey && it.createdAt < vanish.createdAt }) covering[vanish.id] = vanish + if (events.any { it.pubKey == vanish.pubKey && it.createdAt < vanish.createdAt }) covering[vanish.id] = vanish } return covering.values.toList() } + +/** [deletionsCovering] with the local store as the deletion source (the "up" direction). */ +suspend fun IEventStore.deletionsCovering( + serverEvents: List, + relay: NormalizedRelayUrl, +): List = deletionsCovering(serverEvents, relay) { query(it) }