mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
refactor: extract deletion-settle loop into a quartz INostrClient accessory
The two-pass deletion convergence is protocol logic, not CLI assembly, and the geode mirror is a near-term second consumer — so move it out of SyncCommand into a reusable accessory alongside the rest of the negentropy family. quartz: negentropySettleDeletions(relay, filter, store, sendUp, applyDown, …) — re-reconciles after a content settle and resolves only the residual: publishes our covering deletions up (sendUp) and/or ingests the relay's kind-5 down (applyDown, vanish never auto-applied), looping until a round resolves nothing. Returns DeletionSettleResult(sentUp, appliedDown, rounds). Everything it needs is already quartz (negentropyReconcileIds, fetchAll, deletionsCovering, publishAndConfirm, Event.verify, IEventStore), so it carries no CLI dependency. SyncCommand's pass 2 collapses to a single call; pass 1 (content) is unchanged. Catalogued in the accessories README. Tests: DeletionSyncTest drives the accessory end-to-end both ways (sendUp → relay converges to gone; applyDown → local converges to gone), on top of the existing deletionsCovering unit + manual-wiring cases. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JgL1WTV4Hkp2uuXcUHCHGt
This commit is contained in:
@@ -26,15 +26,13 @@ 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.relay.client.accessories.DeletionSettleResult
|
||||
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.client.accessories.negentropySettleDeletions
|
||||
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
|
||||
@@ -186,74 +184,27 @@ 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.
|
||||
// ── Pass 2+: deletion settle. The reusable quartz accessory re-reconciles
|
||||
// and resolves only the residual — send our deletions up for what we deleted
|
||||
// (bounded by --down), apply the relay's kind-5 down for what it deleted
|
||||
// (bounded by --up) — looping until stable. Cheap regardless of database size
|
||||
// (see negentropySettleDeletions), and best-effort so it can't fail the sync.
|
||||
val deletions =
|
||||
if (syncDeletions) {
|
||||
ctx.client.negentropySettleDeletions(
|
||||
relay = relay,
|
||||
filter = filter,
|
||||
store = ctx.store,
|
||||
sendUp = down,
|
||||
applyDown = up,
|
||||
batchSize = ID_CHUNK,
|
||||
idleTimeoutMs = timeoutMs,
|
||||
maxRounds = MAX_DELETION_ROUNDS,
|
||||
reconcileConcurrency = RECONCILE_CONCURRENCY,
|
||||
)
|
||||
} else {
|
||||
DeletionSettleResult(0, 0, 0)
|
||||
}
|
||||
}
|
||||
|
||||
Output.emit(
|
||||
mapOf(
|
||||
@@ -264,9 +215,9 @@ object SyncCommand {
|
||||
"have" to result.haveCount,
|
||||
"downloaded" to downloaded.get(),
|
||||
"uploaded" to uploaded.get(),
|
||||
"deletions_sent_up" to deletionsUp,
|
||||
"deletions_applied_down" to deletionsDown,
|
||||
"deletion_rounds" to deletionRounds,
|
||||
"deletions_sent_up" to deletions.sentUp,
|
||||
"deletions_applied_down" to deletions.appliedDown,
|
||||
"deletion_rounds" to deletions.rounds,
|
||||
),
|
||||
)
|
||||
return 0
|
||||
|
||||
@@ -26,6 +26,7 @@ import com.vitorpamplona.geode.testing.publish
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.negentropyReconcileIds
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.negentropySettleDeletions
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
|
||||
@@ -196,4 +197,71 @@ class DeletionSyncTest : RelayClientTest() {
|
||||
"local applied the pulled deletion and removed the note",
|
||||
)
|
||||
}
|
||||
|
||||
// ---- the full accessory loop (negentropySettleDeletions) -----------------
|
||||
|
||||
// sendUp: local holds the deletion, relay still has the note → the loop pushes it
|
||||
// up and the relay converges to gone.
|
||||
@Test
|
||||
fun settleSendsOurDeletionUp() =
|
||||
runBlocking {
|
||||
val target = note("settle up")
|
||||
val deletion = signer.sign(DeletionEvent.build(listOf(target), createdAt = target.createdAt + 1))
|
||||
val localStore = EventStore(null)
|
||||
localStore.insert(target)
|
||||
localStore.insert(deletion) // deletes target locally, keeps the kind-5
|
||||
defaultRelay.preload(listOf(target))
|
||||
|
||||
val res =
|
||||
withTimeout(30_000) {
|
||||
client.negentropySettleDeletions(
|
||||
relay = defaultRelayUrl,
|
||||
filter = Filter(kinds = listOf(1)),
|
||||
store = localStore,
|
||||
sendUp = true,
|
||||
applyDown = false,
|
||||
idleTimeoutMs = 20_000,
|
||||
)
|
||||
}
|
||||
|
||||
assertEquals(1, res.sentUp)
|
||||
assertEquals(0, res.appliedDown)
|
||||
assertTrue(
|
||||
defaultRelay.store.query<Event>(Filter(ids = listOf(target.id))).isEmpty(),
|
||||
"relay converged: the deleted note is gone",
|
||||
)
|
||||
localStore.close()
|
||||
}
|
||||
|
||||
// applyDown: relay deleted the note (holds only the kind-5), local still has it →
|
||||
// the loop pulls the relay's deletion down and local converges to gone.
|
||||
@Test
|
||||
fun settleAppliesRelayDeletionDown() =
|
||||
runBlocking {
|
||||
val target = note("settle down")
|
||||
val deletion = signer.sign(DeletionEvent.build(listOf(target), createdAt = target.createdAt + 1))
|
||||
defaultRelay.preload(listOf(target, deletion)) // relay deletes target, keeps the kind-5
|
||||
val localStore = EventStore(null)
|
||||
localStore.insert(target)
|
||||
|
||||
val res =
|
||||
withTimeout(30_000) {
|
||||
client.negentropySettleDeletions(
|
||||
relay = defaultRelayUrl,
|
||||
filter = Filter(kinds = listOf(1)),
|
||||
store = localStore,
|
||||
sendUp = false,
|
||||
applyDown = true,
|
||||
idleTimeoutMs = 20_000,
|
||||
)
|
||||
}
|
||||
|
||||
assertEquals(0, res.sentUp)
|
||||
assertEquals(1, res.appliedDown)
|
||||
assertTrue(
|
||||
localStore.query<Event>(Filter(ids = listOf(target.id))).isEmpty(),
|
||||
"local converged: the relay-deleted note is gone",
|
||||
)
|
||||
localStore.close()
|
||||
}
|
||||
}
|
||||
|
||||
+145
@@ -0,0 +1,145 @@
|
||||
/*
|
||||
* 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.quartz.nip01Core.relay.client.accessories
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.crypto.verify
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.store.IEventStore
|
||||
import com.vitorpamplona.quartz.nip01Core.store.deletionsCovering
|
||||
import com.vitorpamplona.quartz.nip09Deletions.DeletionEvent
|
||||
|
||||
/**
|
||||
* Outcome of a [negentropySettleDeletions] run.
|
||||
*
|
||||
* @property sentUp distinct local deletions published to the relay (up direction).
|
||||
* @property appliedDown distinct relay deletions ingested into [store] (down direction).
|
||||
* @property rounds reconcile rounds run before convergence (or the cap).
|
||||
*/
|
||||
class DeletionSettleResult(
|
||||
val sentUp: Int,
|
||||
val appliedDown: Int,
|
||||
val rounds: Int,
|
||||
)
|
||||
|
||||
/**
|
||||
* Converge deletions between [store] and [relay] AFTER a content sync has settled the
|
||||
* two sides — the second half of a two-pass sync. NIP-77 reconciles by id, so a plain
|
||||
* content sync converges everything except events a deletion physically stops from
|
||||
* moving; those survive as the reconcile's residual, which this resolves:
|
||||
*
|
||||
* - **[sendUp]** — a residual **need** (relay has it, [store] still lacks it after the
|
||||
* content pass tried to download it) means we deleted it. Publish OUR covering
|
||||
* deletion up ([IEventStore.deletionsCovering]) so the relay drops it.
|
||||
* - **[applyDown]** — a residual **have** ([store] has it, relay still lacks it after
|
||||
* the content pass tried to upload it) means the relay deleted it. Pull the RELAY'S
|
||||
* covering **kind-5** down and ingest it, so [store] drops it too. A NIP-62 vanish is
|
||||
* deliberately NOT applied on pull — its blast radius is the author's whole account.
|
||||
*
|
||||
* Because it works off the residual — not every id — the cost is one cheap reconcile
|
||||
* per round plus the (small) residual, independent of database size. It loops until a
|
||||
* round resolves nothing (converged, and thereby self-verified) or [maxRounds] is hit.
|
||||
*
|
||||
* **Direction requires the matching content pass.** A residual need is a clean signal
|
||||
* only after the content sync attempted the download ([sendUp] pairs with a `--down`
|
||||
* content pass); a residual have only after it attempted the upload ([applyDown] pairs
|
||||
* with `--up`). Passing a direction whose content pass didn't run makes its residual the
|
||||
* full unsettled set, not a deletion signal — so drive this with the same directions the
|
||||
* content pass used.
|
||||
*
|
||||
* Best-effort: a reconcile failure ([NegentropySyncException]) stops the loop and returns
|
||||
* what already settled rather than throwing — the content sync is the primary work.
|
||||
*
|
||||
* @param batchSize ids per reconcile chunk and per by-id fetch.
|
||||
* @param idleTimeoutMs idle watchdog for the reconciles and fetches.
|
||||
* @param maxRounds hard cap on rounds; the "resolved nothing" check usually stops first.
|
||||
* @param reconcileConcurrency overlapped `created_at`-window reconciles after an over-cap split.
|
||||
*/
|
||||
suspend fun INostrClient.negentropySettleDeletions(
|
||||
relay: NormalizedRelayUrl,
|
||||
filter: Filter,
|
||||
store: IEventStore,
|
||||
sendUp: Boolean,
|
||||
applyDown: Boolean,
|
||||
batchSize: Int = 500,
|
||||
idleTimeoutMs: Long = 120_000L,
|
||||
maxRounds: Int = 4,
|
||||
reconcileConcurrency: Int = 1,
|
||||
): DeletionSettleResult {
|
||||
if ((!sendUp && !applyDown) || maxRounds <= 0) return DeletionSettleResult(0, 0, 0)
|
||||
|
||||
val publishTimeoutSecs = (idleTimeoutMs / 1000).coerceAtLeast(1)
|
||||
val sentUp = HashSet<HexKey>()
|
||||
val appliedDown = HashSet<HexKey>()
|
||||
var rounds = 0
|
||||
|
||||
while (rounds < maxRounds) {
|
||||
rounds++
|
||||
val diff =
|
||||
try {
|
||||
negentropyReconcileIds(
|
||||
relay = relay,
|
||||
filter = filter,
|
||||
localEntries = store.snapshotIdsForNegentropy(listOf(filter)),
|
||||
batchSize = batchSize,
|
||||
idleTimeoutMs = idleTimeoutMs,
|
||||
reconcileConcurrency = reconcileConcurrency,
|
||||
)
|
||||
} catch (e: NegentropySyncException) {
|
||||
break
|
||||
}
|
||||
|
||||
var resolved = 0
|
||||
|
||||
// residual needs → publish our covering deletions up.
|
||||
if (sendUp) {
|
||||
for (chunk in diff.needIds.chunked(batchSize)) {
|
||||
val events = fetchAll(relay, Filter(ids = chunk), idleTimeoutMs)
|
||||
for (del in store.deletionsCovering(events, relay)) {
|
||||
if (sentUp.add(del.id)) {
|
||||
if (publishAndConfirm(del, setOf(relay), publishTimeoutSecs)) resolved++
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// residual haves → ingest the relay's covering kind-5 (never a vanish).
|
||||
if (applyDown) {
|
||||
for (chunk in diff.haveIds.chunked(batchSize)) {
|
||||
val ours = store.query<Event>(Filter(ids = chunk))
|
||||
val relayDeletions = deletionsCovering(ours, relay) { f -> fetchAll(relay, f, idleTimeoutMs) }
|
||||
for (del in relayDeletions.filterIsInstance<DeletionEvent>()) {
|
||||
if (del.verify() && appliedDown.add(del.id)) {
|
||||
store.insert(del)
|
||||
resolved++
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (resolved == 0) break
|
||||
}
|
||||
|
||||
return DeletionSettleResult(sentUp.size, appliedDown.size, rounds)
|
||||
}
|
||||
+1
@@ -52,6 +52,7 @@ Import as `com.vitorpamplona.quartz.nip01Core.relay.client.accessories.<name>` (
|
||||
| `negentropySyncEvents` / `negentropySyncOrFetchEvents` | `NostrClientNegentropySyncEventsExt` | The two above as an O(1)-memory `Flow<Event>`. |
|
||||
| `negentropyReconcile(relay, filter, localEntries, onNeedIds, onHaveIds)` | `NostrClientNegentropySyncExt` | **Pure diff, no I/O** — streams the two directions (`need` = relay has & we lack; `have` = we have & relay lacks) to callbacks. Compose your own download/upload on top. |
|
||||
| `negentropyReconcileIds(relay, filter, localEntries)` | `NostrClientNegentropySyncExt` | Same diff, materialized into `needIds` / `haveIds` lists (small sets only). |
|
||||
| `negentropySettleDeletions(relay, filter, store, sendUp, applyDown)` | `NostrClientNegentropyDeletionSettleExt` | Second pass of a two-pass sync: after a content sync settles, re-reconcile and resolve only the residual — send our covering deletions up (`sendUp`) and/or apply the relay's kind-5 down (`applyDown`), looping until stable. Cost is O(residual), not O(db). Pairs with `IEventStore.deletionsCovering`. |
|
||||
|
||||
`fetchByIds`, `reconcileStreaming`, `syncPipeline` in `NostrClientNegentropySyncExt`
|
||||
are `internal` implementation details — not part of the public surface.
|
||||
|
||||
Reference in New Issue
Block a user