feat(quartz): make FTS reindex pausable/resumable for large stores

A full FTS rebuild can run for a long time on a big store, so add a
resumable, batched overload alongside the one-shot:

    reindexFullTextSearch(resumeFrom: String?, batchSize): FtsReindexProgress

Each call processes ~batchSize events in its own write transaction and
returns an opaque cursor + done flag. The caller loops until done and may
stop at any point — the cursor is durable across crash/app-restart, and
the writer lock is released between batches, so "pause" is just "don't
make the next call". The path is additive/refresh and keeps search usable
throughout (no up-front wipe); the one-shot variant remains for a
guaranteed-clean rebuild.

- SQLite: FullTextSearchModule.reindexBatch walks event_headers ordered
  by the monotonic row_id (a free, stable cursor), restricted to
  searchable kinds, delete-then-insert per event so batches are
  idempotent and never duplicate rows.
- Filesystem: FsEventStore walks one idx/kind/<k>/ dir per step (linear,
  no re-sort); cursor is the next searchable kind. Idempotent linkFts, so
  nothing is wiped. Pauses between kinds.
- Wrappers delegate; new FtsReindexProgress value type carries cursor +
  progress + done.
- cli: `amy store reindex-fts` now loops the batched path to completion
  and reports processed/batch counts.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01BZqPFds2TPPUKkMmBngwys
This commit is contained in:
Claude
2026-06-18 22:16:48 +00:00
parent 0c67f63bad
commit b419db4fe0
11 changed files with 390 additions and 1 deletions
@@ -23,6 +23,7 @@ package com.vitorpamplona.amethyst.cli.commands
import com.vitorpamplona.amethyst.cli.DataDir
import com.vitorpamplona.amethyst.cli.Output
import com.vitorpamplona.quartz.nip01Core.jackson.JacksonMapper
import com.vitorpamplona.quartz.nip01Core.store.IEventStore
import com.vitorpamplona.quartz.nip01Core.store.fs.FsEventStore
import java.io.IOException
import java.nio.file.Files
@@ -169,11 +170,25 @@ object StoreCommands {
withStore(dataDir) { store ->
val ftsDir = dataDir.eventsDir.toPath().resolve("idx/fts")
val before = countEntries(ftsDir)
store.reindexFullTextSearch()
// Drive the resumable, batched path to completion so a huge
// store is processed without holding the writer lock for the
// whole pass. A real long-running caller would persist the
// cursor between calls; here we just loop until done.
var cursor: String? = null
var processed = 0L
var batches = 0
do {
val progress = store.reindexFullTextSearch(cursor, IEventStore.DEFAULT_FTS_REINDEX_BATCH)
cursor = progress.cursor
processed += progress.processedThisBatch
batches++
} while (!progress.done)
val after = countEntries(ftsDir)
Output.emit(
mapOf(
"ok" to true,
"processed" to processed,
"batches" to batches,
"tokens_before" to before,
"tokens_after" to after,
),
@@ -23,6 +23,7 @@ package com.vitorpamplona.quartz.nip01Core.cache.interning
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.store.FtsReindexProgress
import com.vitorpamplona.quartz.nip01Core.store.IEventStore
import com.vitorpamplona.quartz.nip01Core.store.IdAndTime
@@ -125,5 +126,10 @@ class InterningEventStore(
override suspend fun reindexFullTextSearch() = inner.reindexFullTextSearch()
override suspend fun reindexFullTextSearch(
resumeFrom: String?,
batchSize: Int,
): FtsReindexProgress = inner.reindexFullTextSearch(resumeFrom, batchSize)
override fun close() = inner.close()
}
@@ -0,0 +1,42 @@
/*
* 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.store
/**
* Progress token returned by the resumable
* [IEventStore.reindexFullTextSearch] overload.
*
* Treat [cursor] as opaque: persist it (it survives process death) and
* hand it back to the next call to continue where the previous one
* stopped. Each store encodes its own resume position into it.
*
* @property cursor where to resume from on the next call, or `null` once
* [done] is `true` (nothing left to process).
* @property processedThisBatch how many events this call (re)indexed —
* useful to drive a progress indicator.
* @property done `true` when the whole store has been visited; further
* calls are no-ops that keep returning `done = true`.
*/
data class FtsReindexProgress(
val cursor: String?,
val processedThisBatch: Int,
val done: Boolean,
)
@@ -25,6 +25,16 @@ import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
interface IEventStore : AutoCloseable {
companion object {
/**
* How many events a single resumable
* [reindexFullTextSearch] batch aims to process before yielding.
* Big enough to amortise the per-batch transaction/lock cost,
* small enough that a pause request is honoured promptly.
*/
const val DEFAULT_FTS_REINDEX_BATCH = 1000
}
/**
* Relay URL this store is acting on behalf of, or `null` for an
* unscoped store. Used by NIP-62 right-to-vanish handling: only
@@ -153,8 +163,48 @@ interface IEventStore : AutoCloseable {
* that currently map to a searchable event, leaving the bulk of
* non-searchable rows (reactions, zaps, follow lists, …) untouched
* to keep the scan as cheap as possible.
*
* This one-shot variant runs to completion under a single lock and
* cannot be paused. For a store large enough that the full pass
* would block for too long, use the resumable overload
* [reindexFullTextSearch] instead.
*/
suspend fun reindexFullTextSearch()
/**
* Resumable, batched companion to [reindexFullTextSearch] for stores
* large enough that a single pass would take too long to run (or to
* hold a lock) in one go.
*
* Process roughly [batchSize] events starting from [resumeFrom]
* (`null` = from the beginning) and return a [FtsReindexProgress].
* Drive it in a loop, feeding [FtsReindexProgress.cursor] back in,
* until [FtsReindexProgress.done] is `true`:
*
* ```
* var cursor: String? = null
* do {
* val p = store.reindexFullTextSearch(cursor)
* cursor = p.cursor
* // optionally persist `cursor` and stop; resume later by passing it back
* } while (!p.done)
* ```
*
* Each call commits its own batch, so progress is durable across a
* crash or app restart and the writer lock is released between
* batches — the app can pause simply by not making the next call.
*
* Semantics differ slightly from the one-shot variant: this path is
* **additive / refresh** — it makes sure every currently-searchable
* event is indexed (and, on SQLite, refreshes changed content) while
* leaving search usable throughout. It does not purge stale rows left
* behind by a kind that *lost* searchability; run the one-shot
* [reindexFullTextSearch] once for that rarer case.
*/
suspend fun reindexFullTextSearch(
resumeFrom: String?,
batchSize: Int = DEFAULT_FTS_REINDEX_BATCH,
): FtsReindexProgress
override fun close()
}
@@ -220,5 +220,10 @@ class ObservableEventStore(
// projections to observe and no [StoreChange] is emitted.
override suspend fun reindexFullTextSearch() = inner.reindexFullTextSearch()
override suspend fun reindexFullTextSearch(
resumeFrom: String?,
batchSize: Int,
): FtsReindexProgress = inner.reindexFullTextSearch(resumeFrom, batchSize)
override fun close() = inner.close()
}
@@ -25,6 +25,7 @@ import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl
import com.vitorpamplona.quartz.nip01Core.store.FtsReindexProgress
import com.vitorpamplona.quartz.nip01Core.store.IEventStore
import com.vitorpamplona.quartz.nip01Core.store.IdAndTime
@@ -81,5 +82,10 @@ class EventStore(
override suspend fun reindexFullTextSearch() = store.reindexFullTextSearch()
override suspend fun reindexFullTextSearch(
resumeFrom: String?,
batchSize: Int,
): FtsReindexProgress = store.reindexFullTextSearch(resumeFrom, batchSize)
override fun close() = store.close()
}
@@ -24,6 +24,7 @@ import androidx.sqlite.SQLiteConnection
import androidx.sqlite.SQLiteException
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper
import com.vitorpamplona.quartz.nip01Core.store.FtsReindexProgress
import com.vitorpamplona.quartz.nip50Search.SearchableEvent
import com.vitorpamplona.quartz.utils.EventFactory
@@ -77,6 +78,11 @@ class FullTextSearchModule : IModule {
VALUES (?, ?)
""".trimIndent()
val deleteFTSByRowId =
"""
DELETE FROM $tableName WHERE $eventHeaderRowIdName = ?
""".trimIndent()
fun insert(
event: Event,
headerId: Long,
@@ -173,6 +179,81 @@ class FullTextSearchModule : IModule {
}
}
/**
* Process one batch of a resumable rebuild: re-derive the FTS rows
* for up to [batchSize] events whose `row_id > ` [afterRowId] and
* whose kind is searchable, ordered by `row_id`.
*
* `row_id` is a monotonic AUTOINCREMENT key, so it is a stable
* cursor that needs no extra bookkeeping. Each event is
* delete-then-insert, which keeps the batch idempotent (a replay
* after a crash is harmless) and avoids duplicate FTS rows for
* events that were already indexed by the normal insert path. Rows
* not yet reached keep their previous FTS content, so search stays
* usable while the rebuild is in flight.
*
* Must run inside the caller's per-batch write transaction.
*/
fun reindexBatch(
db: SQLiteConnection,
afterRowId: Long,
batchSize: Int,
): FtsReindexProgress {
val kinds = searchableKindsPresent(db)
if (kinds.isEmpty()) return FtsReindexProgress(cursor = null, processedThisBatch = 0, done = true)
val selectSql =
"SELECT row_id, id, pubkey, created_at, kind, tags, content, sig " +
"FROM event_headers WHERE row_id > ? AND kind IN (${kinds.joinToString(",")}) " +
"ORDER BY row_id LIMIT ?"
var last = afterRowId
var processed = 0
db.prepare(deleteFTSByRowId).use { del ->
db.prepare(insertFTS).use { write ->
db.prepare(selectSql).use { read ->
read.bindLong(1, afterRowId)
read.bindLong(2, batchSize.toLong())
while (read.step()) {
val rowId = read.getLong(0)
// Clear any existing row for this event first so a
// replay or an already-indexed event can't duplicate.
del.bindLong(1, rowId)
del.step()
del.reset()
val event =
EventFactory.create<Event>(
read.getText(1),
read.getText(2),
read.getLong(3),
read.getInt(4),
OptimizedJsonMapper.fromJsonToTagArray(read.getText(5)),
read.getText(6),
read.getText(7),
)
if (event is SearchableEvent) {
write.bindLong(1, rowId)
write.bindText(2, event.indexableContent())
write.step()
write.reset()
}
last = rowId
processed++
}
}
}
}
// Fewer than a full page came back ⇒ we hit the end of the table.
val done = processed < batchSize
return FtsReindexProgress(
cursor = if (done) null else last.toString(),
processedThisBatch = processed,
done = done,
)
}
/**
* The distinct kinds present in `event_headers` that currently parse
* to a [SearchableEvent]. Kind alone selects the event class in
@@ -31,6 +31,7 @@ import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper
import com.vitorpamplona.quartz.nip01Core.core.isEphemeral
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.store.FtsReindexProgress
import com.vitorpamplona.quartz.nip01Core.store.IEventStore
import com.vitorpamplona.quartz.nip01Core.store.IdAndTime
import com.vitorpamplona.quartz.nip40Expiration.isExpired
@@ -353,6 +354,22 @@ class SQLiteEventStore(
}
}
/**
* One batch of a resumable FTS rebuild. See
* [IEventStore.reindexFullTextSearch]. The opaque cursor is the last
* `row_id` processed; `null` (or an unparseable value) starts from
* the beginning. Each batch is its own write transaction.
*/
suspend fun reindexFullTextSearch(
resumeFrom: String?,
batchSize: Int,
): FtsReindexProgress =
pool.useWriter { db ->
db.transaction {
fullTextSearchModule.reindexBatch(db, resumeFrom?.toLongOrNull() ?: 0L, batchSize)
}
}
fun close() = pool.close()
}
@@ -230,6 +230,75 @@ class SearchTest : BaseDBTest() {
db.store.assertQuery(note, Filter(search = "uniqnote"))
}
@Test
fun testResumableReindexProcessesEveryEventInBatches() =
forEachDB { db ->
// A handful of searchable notes, each with a unique token.
val notes =
(0 until 5).map { i ->
signer.sign(TextNoteEvent.build("uniqresume$i body", createdAt = TimeUtils.now() + i))
}
notes.forEach { db.store.insertEvent(it) }
// Mimic the post-upgrade state: canonical rows present, FTS empty.
db.store.pool.useWriter { db.store.fullTextSearchModule.deleteAll(it) }
db.store.assertQuery(null, Filter(search = "uniqresume0"))
// Drive the resumable path two events at a time, persisting only
// the opaque cursor between calls — exactly what a paused/resumed
// app would do.
var cursor: String? = null
var batches = 0
var processed = 0
do {
val progress = db.store.reindexFullTextSearch(cursor, batchSize = 2)
cursor = progress.cursor
processed += progress.processedThisBatch
batches++
} while (!progress.done)
// 5 events at 2 per batch ⇒ 3 batches (2 + 2 + 1).
kotlin.test.assertEquals(5, processed)
kotlin.test.assertEquals(3, batches)
// Every note is searchable again after the full run.
notes.forEachIndexed { i, n -> db.store.assertQuery(n, Filter(search = "uniqresume$i")) }
}
@Test
fun testResumableReindexKeepsSearchLiveAndIsIdempotent() =
forEachDB { db ->
val a = signer.sign(TextNoteEvent.build("uniqalpha note", createdAt = TimeUtils.now()))
val b = signer.sign(TextNoteEvent.build("uniqbeta note", createdAt = TimeUtils.now() + 1))
db.store.insertEvent(a)
db.store.insertEvent(b)
// First batch (size 1) reindexes only the lowest row_id; the
// untouched event keeps the FTS row it already had from insert,
// so search stays live for both throughout.
val first = db.store.reindexFullTextSearch(null, batchSize = 1)
kotlin.test.assertEquals(false, first.done)
db.store.assertQuery(a, Filter(search = "uniqalpha"))
db.store.assertQuery(b, Filter(search = "uniqbeta"))
// Finish, then run the whole loop again from scratch: delete-then-
// insert per event means no duplicate rows (assertQuery wants 1).
var cursor = first.cursor
do {
val p = db.store.reindexFullTextSearch(cursor, batchSize = 1)
cursor = p.cursor
} while (!p.done)
cursor = null
do {
val p = db.store.reindexFullTextSearch(cursor, batchSize = 10)
cursor = p.cursor
} while (!p.done)
db.store.assertQuery(a, Filter(search = "uniqalpha"))
db.store.assertQuery(b, Filter(search = "uniqbeta"))
}
@Test
fun testChannelJsonFieldsAreSearchable() =
forEachDB { db ->
@@ -28,6 +28,7 @@ import com.vitorpamplona.quartz.nip01Core.core.isEphemeral
import com.vitorpamplona.quartz.nip01Core.core.isReplaceable
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.store.FtsReindexProgress
import com.vitorpamplona.quartz.nip01Core.store.IEventStore
import com.vitorpamplona.quartz.nip01Core.store.sqlite.DefaultIndexingStrategy
import com.vitorpamplona.quartz.nip01Core.store.sqlite.IndexingStrategy
@@ -528,6 +529,63 @@ open class FsEventStore(
}
}
/**
* Resumable, additive FTS rebuild. See
* [IEventStore.reindexFullTextSearch] for the loop contract.
*
* The cursor is the next searchable kind to process, so the store is
* walked one `idx/kind/<k>/` directory at a time — that keeps the
* scan linear (each directory is listed once, never re-sorted) at the
* cost of pausing only between kinds: a single huge kind is processed
* in one batch. [batchSize] is a soft floor — whole kinds are
* processed until at least that many events have been (re)linked,
* then the call yields. Linking is idempotent ([FsIndexer.linkFts]
* ignores an existing hardlink), so nothing is wiped and search stays
* usable throughout; replaying a kind after a crash is harmless.
*/
override suspend fun reindexFullTextSearch(
resumeFrom: String?,
batchSize: Int,
): FtsReindexProgress =
lockManager.withWriteLock {
if (!Files.isDirectory(layout.idxKind)) {
return@withWriteLock FtsReindexProgress(cursor = null, processedThisBatch = 0, done = true)
}
Files.createDirectories(layout.idxFts)
val resumeKind = resumeFrom?.toIntOrNull()
val pending =
Files
.list(layout.idxKind)
.use { dirs -> dirs.map { it.fileName.toString() }.toList() }
.mapNotNull { it.toIntOrNull() }
.filter { isSearchableKind(it) && (resumeKind == null || it >= resumeKind) }
.sorted()
var processed = 0
var index = 0
while (index < pending.size) {
val kind = pending[index]
Files.list(layout.kindDir(kind)).use { entries ->
for (entry in entries) {
val id = FsLayout.parseEntry(entry.fileName.toString())?.second ?: continue
val event = readEvent(id) ?: continue
indexer.linkFts(event, layout.canonical(id))
processed++
}
}
index++
if (processed >= batchSize) break
}
val done = index >= pending.size
FtsReindexProgress(
cursor = if (done) null else pending[index].toString(),
processedThisBatch = processed,
done = done,
)
}
/**
* True when [kind] currently parses to a [SearchableEvent]. Kind
* alone selects the event class in [EventFactory], so a single probe
@@ -24,6 +24,7 @@ import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync
import com.vitorpamplona.quartz.nip10Notes.TextNoteEvent
import com.vitorpamplona.quartz.nip23LongContent.LongTextNoteEvent
import com.vitorpamplona.quartz.utils.Secp256k1Instance
import kotlinx.coroutines.runBlocking
import java.nio.file.Files
@@ -159,6 +160,45 @@ class FsSearchTest {
assertEquals(2, ftsRoot.resolve("reindex").listDirectoryEntries().size)
}
@Test
fun `resumable reindex covers every kind across batches`() =
runBlocking {
// Two searchable kinds (note = kind 1, long-form = kind 30023) so
// the kind-granular cursor must advance across more than one dir.
val n = note("uniqnote bitcoin", ts = 100)
val long =
signer.sign(
LongTextNoteEvent.build(
"uniqlong body",
title = "title",
dTag = "d1",
createdAt = 200,
),
)
store.insert(n)
store.insert(long)
// Wipe the index, then drive the resumable path one kind at a time.
val ftsRoot = root.resolve("idx/fts")
Files.walk(ftsRoot).use { stream ->
stream.sorted(Comparator.reverseOrder()).forEach { p -> if (p != ftsRoot) Files.deleteIfExists(p) }
}
assertTrue(store.query<Event>(Filter(search = "uniqnote")).isEmpty())
var cursor: String? = null
var batches = 0
do {
val progress = store.reindexFullTextSearch(cursor, batchSize = 1)
cursor = progress.cursor
batches++
} while (!progress.done)
// Two searchable kind dirs, batchSize 1 ⇒ at least two batches.
assertTrue(batches >= 2, "expected the cursor to span both kinds, got $batches batch(es)")
assertEquals(listOf(n.id), store.query<Event>(Filter(search = "uniqnote")).map { it.id })
assertEquals(listOf(long.id), store.query<Event>(Filter(search = "uniqlong")).map { it.id })
}
@Test
fun `reindexFullTextSearch ignores non-searchable kinds`() =
runBlocking {