mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
refactor: deletion sync as a post-settle residual pass (both directions, O(residual))
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 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JgL1WTV4Hkp2uuXcUHCHGt
This commit is contained in:
@@ -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<String>,
|
||||
@@ -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<HexKey>()
|
||||
|
||||
// ── 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<List<HexKey>>(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<List<HexKey>>(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<HexKey>()
|
||||
val appliedDown = HashSet<HexKey>()
|
||||
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<Event>(Filter(ids = chunk))
|
||||
val relayDeletions = deletionsCovering(ourEvents, relay) { f -> ctx.client.fetchAll(relay, f, timeoutMs) }
|
||||
for (del in relayDeletions.filterIsInstance<DeletionEvent>()) {
|
||||
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
|
||||
|
||||
@@ -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<Event>(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<Event>(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<Event>(Filter(ids = diff.haveIds))
|
||||
val relayDeletions = deletionsCovering(ourEvents, defaultRelayUrl) { f -> defaultRelay.store.query<Event>(f) }
|
||||
assertEquals(listOf(deletion.id), relayDeletions.map { it.id })
|
||||
relayDeletions.filterIsInstance<DeletionEvent>().forEach { local.store.insert(it) }
|
||||
|
||||
assertTrue(
|
||||
local.store.query<Event>(Filter(ids = listOf(target.id))).isEmpty(),
|
||||
"local applied the pulled deletion and removed the note",
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
+26
-12
@@ -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<Event>,
|
||||
suspend fun deletionsCovering(
|
||||
events: List<Event>,
|
||||
relay: NormalizedRelayUrl,
|
||||
query: suspend (Filter) -> List<Event>,
|
||||
): List<Event> {
|
||||
if (serverEvents.isEmpty()) return emptyList()
|
||||
if (events.isEmpty()) return emptyList()
|
||||
val covering = LinkedHashMap<HexKey, Event>()
|
||||
|
||||
// 1. id-based NIP-09: a kind-5 `e`-tagging a server id.
|
||||
query<Event>(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<Event>(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<Event>(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<Event>,
|
||||
relay: NormalizedRelayUrl,
|
||||
): List<Event> = deletionsCovering(serverEvents, relay) { query<Event>(it) }
|
||||
|
||||
Reference in New Issue
Block a user