From 6ac2903a3a52132984711ecab8eb4a999f22fde1 Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 4 Jul 2026 02:31:40 +0000 Subject: [PATCH] =?UTF-8?q?fix(quartz):=20live=20negentropy=20index=20?= =?UTF-8?q?=E2=80=94=20same-batch=20displacement,=20no-op=20kind-5,=20rebu?= =?UTF-8?q?ild=20cap?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Audit follow-ups on the live NIP-77 index, all with the store/scan equivalence test extended to cover them: - Same-batch replaceable displacement left a dead id in the index. applyAfterCommit applied all removes before all adds, so when a later row in one transaction displaced an earlier row of the same batch (two versions of one replaceable — the mirror-backfill hot path), the displaced row's remove no-op'd against an index that hadn't taken its add yet, then the add re-inserted it: the index advertised an id the trigger had already deleted. recordAccepted now cancels the pending add instead of queueing a remove (added is a LinkedHashSet for O(1) cancel). - A kind-5 that deleted nothing (a delete broadcast for events this relay never stored — the common case) still invalidated the whole index, forcing a full-scan rebuild under the writer mutex on the next NEG-OPEN. DeletionRequestModule.insert now returns the rows it deleted; recordAccepted only invalidates when that count is > 0, else records the kind-5 as a plain row. - The first NEG-OPEN over a corpus larger than the serve cap scanned the whole table uncapped, built a full index that could never produce a snapshot, and then maintained it forever for zero benefit. liveNegentropySnapshot now caps the rebuild scan at maxEntries + 1 and leaves the index unpopulated when the corpus is over-cap (the scan path answers NEG-ERR, as before). - delete/deleteExpired/clearDB invalidate() moved inside the writer mutex so no NEG-OPEN can seal a snapshot of just-deleted rows, and no concurrent rebuild can be discarded by a late invalidate. Also corrects the ~40 B/event heap figure to ~140 B (IdAndTime keeps the id as a 64-char hex string, not 32 bytes) in the strategy/plan docs. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01TtDNpayEYvJH7QuPswND3A --- .../geode/RelayIndexingStrategy.kt | 4 +- ...26-07-03-incremental-negentropy-storage.md | 6 +- .../store/sqlite/DeletionRequestModule.kt | 17 ++- .../store/sqlite/IndexingStrategy.kt | 5 +- .../store/sqlite/SQLiteEventStore.kt | 94 +++++++++---- .../sqlite/LiveNegentropyIndexStoreTest.kt | 130 ++++++++++++++++++ 6 files changed, 214 insertions(+), 42 deletions(-) diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/RelayIndexingStrategy.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/RelayIndexingStrategy.kt index 9b3a116f3f..33fa9f4441 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/RelayIndexingStrategy.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/RelayIndexingStrategy.kt @@ -43,8 +43,8 @@ import com.vitorpamplona.quartz.nip01Core.store.sqlite.DefaultIndexingStrategy * @param liveNegentropyIndex keep the always-current `(created_at, id)` * set that serves full-corpus NIP-77 NEG-OPENs without a scan + seal * (strfry answers those off its live tree; the scan+seal path measured - * ~340 ms per cold open at 50k events). Costs ~40 B/event of heap and - * one indexed pre-SELECT per replaceable insert. + * ~340 ms per cold open at 50k events). Costs ~140 B/event of JVM heap + * (hex-string ids) and one indexed pre-SELECT per replaceable insert. * `[negentropy].live_index = false` turns it off. */ fun relayIndexingStrategy( diff --git a/quartz/plans/2026-07-03-incremental-negentropy-storage.md b/quartz/plans/2026-07-03-incremental-negentropy-storage.md index ca10f6b091..d3e7ffcc80 100644 --- a/quartz/plans/2026-07-03-incremental-negentropy-storage.md +++ b/quartz/plans/2026-07-03-incremental-negentropy-storage.md @@ -60,8 +60,10 @@ are effectively always cold; any new filter is cold by definition. ### 1. `LiveNegentropyIndex` (quartz, server-only, opt-in) An always-current sorted set of `(createdAt, id₃₂)` maintained from the -store's write path — strfry's `MemoryView` equivalent, ~40 B/entry -(1M events ≈ 40 MB; capped by `negentropy.max_sync_events`). +store's write path — strfry's `MemoryView` equivalent, ~140 B/entry on +the JVM (`IdAndTime` keeps the id as a 64-char hex string; 1M events ≈ +140 MB; the index is only built when the corpus fits +`negentropy.max_sync_events`, so that also caps the heap). - **Structure**: single sorted array with binary-search insert. Nostr inserts are near-tail (created_at ≈ now), so the memmove is tiny in the common diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/DeletionRequestModule.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/DeletionRequestModule.kt index 30bb25f99e..448296a4ae 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/DeletionRequestModule.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/DeletionRequestModule.kt @@ -82,15 +82,18 @@ class DeletionRequestModule( fun insert( event: Event, db: SQLiteConnection, - ) { - if (event is DeletionEvent) { - val idValues = event.deleteEventIds() - val addresses = event.deleteAddresses() + ): Int { + if (event !is DeletionEvent) return 0 - deleteSQL(event.pubKey, event.createdAt, idValues, addresses, hasher(db)).forEach { delete -> - db.execSQL(delete.sql, delete.args) - } + val idValues = event.deleteEventIds() + val addresses = event.deleteAddresses() + + var deletedRows = 0 + deleteSQL(event.pubKey, event.createdAt, idValues, addresses, hasher(db)).forEach { delete -> + db.execSQL(delete.sql, delete.args) + deletedRows += db.changes() } + return deletedRows } /** diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/IndexingStrategy.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/IndexingStrategy.kt index aa7d38ca25..816f50f0f5 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/IndexingStrategy.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/IndexingStrategy.kt @@ -123,8 +123,9 @@ interface IndexingStrategy { * Maintain an always-current in-memory `(created_at, id)` set (a * [com.vitorpamplona.quartz.nip77Negentropy.LiveNegentropyIndex]) so * NIP-77 NEG-OPENs over the full corpus skip the scan + O(n log n) - * seal — strfry answers those off its live tree. Costs ~40 B/event - * of heap plus one indexed pre-SELECT per replaceable insert, which + * seal — strfry answers those off its live tree. Costs ~140 B/event + * of JVM heap (the id is kept as a 64-char hex string, not 32 bytes) + * plus one indexed pre-SELECT per replaceable insert, which * only makes sense on a *relay*; client-side stores don't serve * NEG-OPENs, so the default is **off**. */ diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SQLiteEventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SQLiteEventStore.kt index 94179e9406..8ca31300f7 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SQLiteEventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SQLiteEventStore.kt @@ -214,8 +214,8 @@ class SQLiteEventStore( suspend fun clearDB() { pool.useWriter { db -> modules.reversed().forEach { it.deleteAll(db) } + liveNegentropyIndex?.invalidate() } - liveNegentropyIndex?.invalidate() } suspend fun vacuum() = @@ -278,7 +278,12 @@ class SQLiteEventStore( * when [liveNegentropyIndex] is off: zero cost on the write path. */ internal class LiveIndexDelta { - val added = ArrayList() + /** + * A set (not a list) so a row displaced by a LATER row of the + * same transaction can be cancelled in O(1) — see + * [recordAccepted]. Unique ids make duplicates impossible. + */ + val added = LinkedHashSet() val removed = ArrayList() /** @@ -308,13 +313,28 @@ class SQLiteEventStore( private fun LiveIndexDelta.recordAccepted( event: Event, displaced: IdAndTime?, + deletedRows: Int, ) { - if (event is DeletionEvent || (event is RequestToVanishEvent && event.shouldVanishFrom(relay))) { - // Their own row lands too, but the rebuild scan picks it up. + if (event is RequestToVanishEvent && event.shouldVanishFrom(relay)) { + // Vanish deletes by pubkey inside a trigger; not itemizable. + // Its own row lands too, but the rebuild scan picks it up. invalidateAll = true return } - displaced?.let { removed += it } + if (event is DeletionEvent && deletedRows > 0) { + // Only a kind-5 that actually removed rows costs a rebuild. + // The common case — a delete broadcast for events this relay + // never stored — deletes nothing and records as a plain row, + // so it can't be used to thrash the index. + invalidateAll = true + return + } + // A row displaced by this insert that was added earlier in the + // SAME transaction (two versions of one replaceable in one + // batch) was never in the live index: cancel its pending add. + // Queueing a remove instead would no-op against the index and + // let the stale add through — advertising a dead id. + displaced?.let { if (!added.remove(it)) removed += it } added += IdAndTime(event.createdAt, event.id) } @@ -386,11 +406,11 @@ class SQLiteEventStore( ) { val displaced = if (delta != null) displacedBy(event, db) else null val headerId = eventIndexModule.insert(event, db) - deletionModule.insert(event, db) + val deletedRows = deletionModule.insert(event, db) expirationModule.insert(event, headerId, db) fullTextSearchModule.insert(event, headerId, db) rightToVanishModule.insert(event, relay, headerId, db) - delta?.recordAccepted(event, displaced) + delta?.recordAccepted(event, displaced, deletedRows) } suspend fun insertEvent(event: Event) { @@ -540,35 +560,41 @@ class SQLiteEventStore( maxEntries: Int? = null, ): List = pool.useReader { queryBuilder.snapshotIdsForNegentropy(filters, it, maxEntries) } + // The invalidate() calls below run INSIDE useWriter: dropping the + // index while still holding the mutex means no NEG-OPEN can seal a + // snapshot that still advertises the just-deleted rows, and a + // concurrent rebuild (also writer-mutexed) can't complete only to + // have its fresh content thrown away by a late invalidate. + suspend fun delete(filter: Filter) { - pool.useWriter { queryBuilder.delete(filter, it) } - // Delete-by-filter can't itemize what it removed; drop the live - // index and let the next NEG-OPEN rebuild from one scan. - liveNegentropyIndex?.invalidate() + pool.useWriter { + queryBuilder.delete(filter, it) + // Delete-by-filter can't itemize what it removed; drop the + // live index and let the next NEG-OPEN rebuild from one scan. + liveNegentropyIndex?.invalidate() + } } suspend fun delete(filters: List) { - pool.useWriter { queryBuilder.delete(filters, it) } - liveNegentropyIndex?.invalidate() + pool.useWriter { + queryBuilder.delete(filters, it) + liveNegentropyIndex?.invalidate() + } } - suspend fun delete(id: HexKey): Int { - val changes = - pool.useWriter { db -> - db.execSQL("DELETE FROM event_headers WHERE id = ?", arrayOf(id)) - db.changes() - } - if (changes > 0) liveNegentropyIndex?.invalidate() - return changes - } + suspend fun delete(id: HexKey): Int = + pool.useWriter { db -> + db.execSQL("DELETE FROM event_headers WHERE id = ?", arrayOf(id)) + val changes = db.changes() + if (changes > 0) liveNegentropyIndex?.invalidate() + changes + } suspend fun deleteExpiredEvents() { - val swept = - pool.useWriter { db -> - expirationModule.deleteExpiredEvents(db) - db.changes() - } - if (swept > 0) liveNegentropyIndex?.invalidate() + pool.useWriter { db -> + expirationModule.deleteExpiredEvents(db) + if (db.changes() > 0) liveNegentropyIndex?.invalidate() + } } /** @@ -585,7 +611,17 @@ class SQLiteEventStore( if (!index.isPopulated()) { pool.useWriter { db -> if (!index.isPopulated()) { - index.rebuild(queryBuilder.snapshotIdsForNegentropy(listOf(Filter()), db, null)) + // The scan is capped at maxEntries + 1: a corpus over + // the serve cap can never produce a snapshot (callers + // answer NEG-ERR from the scan path instead), so + // building — and then maintaining — a full index for + // it would cost heap and write-path work for zero + // benefit. Leaving it unpopulated also keeps ingest + // delta-free (see newDeltaOrNull). + val scan = queryBuilder.snapshotIdsForNegentropy(listOf(Filter()), db, maxEntries) + if (scan.size <= maxEntries) { + index.rebuild(scan) + } } } } diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/LiveNegentropyIndexStoreTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/LiveNegentropyIndexStoreTest.kt index 3102f162de..ae04a670df 100644 --- a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/LiveNegentropyIndexStoreTest.kt +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/LiveNegentropyIndexStoreTest.kt @@ -193,6 +193,136 @@ class LiveNegentropyIndexStoreTest { store.close() } + @Test + fun sameBatchReplaceableOverwriteDropsTheDisplacedId() = + runTest { + val store = newStore() + store.insert(event(1)) + // Populate the index first so the batch runs the DELTA path + // (an unpopulated index would be covered by the rebuild scan + // and could never expose this). + assertIndexMatchesScan(store) + + // Three versions of one replaceable in ONE batch: each later + // row displaces the previous one INSIDE the same transaction. + // The displaced rows were never in the index, so their + // pending adds must be cancelled — a remove would no-op and + // leave the index advertising dead ids. + val outcomes = + store.batchInsert( + listOf( + event(2, kind = 0, createdAt = 100, pubKey = pubkey(7)), + event(3, kind = 0, createdAt = 200, pubKey = pubkey(7)), + event(4, kind = 0, createdAt = 300, pubKey = pubkey(7)), + ), + ) + assertTrue(outcomes.all { it == IEventStore.InsertOutcome.Accepted }) + assertEquals(1, store.count(Filter(kinds = listOf(0)))) + assertIndexMatchesScan(store) + + store.close() + } + + @Test + fun sameTransactionReplaceableOverwriteDropsTheDisplacedId() = + runTest { + val store = newStore() + store.insert(event(1)) + assertIndexMatchesScan(store) + + store.transaction { + insert(event(2, kind = 10002, createdAt = 100, pubKey = pubkey(8))) + insert(event(3, kind = 10002, createdAt = 200, pubKey = pubkey(8))) + } + assertEquals(1, store.count(Filter(kinds = listOf(10002)))) + assertIndexMatchesScan(store) + + store.close() + } + + @Test + fun sameBatchAddressableOverwriteDropsTheDisplacedId() = + runTest { + val store = newStore() + store.insert(event(1)) + assertIndexMatchesScan(store) + + val author = pubkey(9) + store.batchInsert( + listOf( + event(2, kind = 30023, createdAt = 100, pubKey = author, tags = arrayOf(arrayOf("d", "post-a"))), + event(3, kind = 30023, createdAt = 200, pubKey = author, tags = arrayOf(arrayOf("d", "post-a"))), + ), + ) + assertEquals(1, store.count(Filter(kinds = listOf(30023)))) + assertIndexMatchesScan(store) + + store.close() + } + + @Test + fun noOpDeletionKeepsTheIndexLive() = + runTest { + val store = newStore() + store.insert(event(1)) + assertIndexMatchesScan(store) + + // A kind-5 referencing an id this store never had — the + // common case for broadcast deletes — removes nothing, so it + // must NOT invalidate the index (that would be a free way to + // force full-scan rebuilds); it lands as a plain row. + store.insert( + event(2, kind = 5, createdAt = 120, tags = arrayOf(arrayOf("e", hexId(999)))), + ) + val index = assertNotNull(store.store.liveNegentropyIndex) + assertTrue(index.isPopulated(), "no-op kind-5 must not invalidate the live index") + assertIndexMatchesScan(store) + + store.close() + } + + @Test + fun effectiveDeletionStillInvalidates() = + runTest { + val store = newStore() + val author = pubkey(3) + val target = event(1, createdAt = 100, pubKey = author) + store.insert(target) + assertIndexMatchesScan(store) + + store.insert( + event(2, kind = 5, createdAt = 120, pubKey = author, tags = arrayOf(arrayOf("e", target.id))), + ) + val index = assertNotNull(store.store.liveNegentropyIndex) + assertTrue(!index.isPopulated(), "a kind-5 that deleted rows must invalidate") + assertIndexMatchesScan(store) + + store.close() + } + + @Test + fun overCapCorpusNeverBuildsTheIndex() = + runTest { + val store = newStore() + store.insert(event(1)) + store.insert(event(2)) + store.insert(event(3)) + + // First-ever NEG-OPEN over a corpus above the serve cap: the + // rebuild scan is capped, the index stays unpopulated (no + // heap, no write-path deltas), and the caller falls back. + assertNull(store.liveNegentropySnapshot(2)) + val index = assertNotNull(store.store.liveNegentropyIndex) + assertTrue(!index.isPopulated(), "an over-cap corpus must not build the index") + + // A cap the corpus fits builds and serves as usual. + assertNotNull(store.liveNegentropySnapshot(3)) + assertTrue(index.isPopulated()) + assertIndexMatchesScan(store) + + store.close() + } + @Test fun overCapFallsBackToNull() = runTest {