mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 11:18:24 +00:00
fix(quartz): live negentropy index — same-batch displacement, no-op kind-5, rebuild cap
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 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TtDNpayEYvJH7QuPswND3A
This commit is contained in:
@@ -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(
|
||||
|
||||
@@ -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
|
||||
|
||||
+10
-7
@@ -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
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
+3
-2
@@ -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**.
|
||||
*/
|
||||
|
||||
+65
-29
@@ -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<IdAndTime>()
|
||||
/**
|
||||
* 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<IdAndTime>()
|
||||
val removed = ArrayList<IdAndTime>()
|
||||
|
||||
/**
|
||||
@@ -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<IdAndTime> = 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<Filter>) {
|
||||
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)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+130
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user