From b419db4fe0a84c85389b585b2a086214e5858ce0 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 18 Jun 2026 22:16:48 +0000 Subject: [PATCH] feat(quartz): make FTS reindex pausable/resumable for large stores MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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// 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 Claude-Session: https://claude.ai/code/session_01BZqPFds2TPPUKkMmBngwys --- .../amethyst/cli/commands/StoreCommands.kt | 17 +++- .../cache/interning/InterningEventStore.kt | 6 ++ .../nip01Core/store/FtsReindexProgress.kt | 42 ++++++++++ .../quartz/nip01Core/store/IEventStore.kt | 50 ++++++++++++ .../nip01Core/store/ObservableEventStore.kt | 5 ++ .../nip01Core/store/sqlite/EventStore.kt | 6 ++ .../store/sqlite/FullTextSearchModule.kt | 81 +++++++++++++++++++ .../store/sqlite/SQLiteEventStore.kt | 17 ++++ .../nip01Core/store/sqlite/SearchTest.kt | 69 ++++++++++++++++ .../quartz/nip01Core/store/fs/FsEventStore.kt | 58 +++++++++++++ .../quartz/nip01Core/store/fs/FsSearchTest.kt | 40 +++++++++ 11 files changed, 390 insertions(+), 1 deletion(-) create mode 100644 quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/FtsReindexProgress.kt diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/StoreCommands.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/StoreCommands.kt index ff9037d6bf..c1ff1b89eb 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/StoreCommands.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/StoreCommands.kt @@ -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, ), diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/cache/interning/InterningEventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/cache/interning/InterningEventStore.kt index 48d1e2faca..2f87070657 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/cache/interning/InterningEventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/cache/interning/InterningEventStore.kt @@ -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() } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/FtsReindexProgress.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/FtsReindexProgress.kt new file mode 100644 index 0000000000..2db5df7cb4 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/FtsReindexProgress.kt @@ -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, +) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/IEventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/IEventStore.kt index 6c7210ba94..7d68d4fb93 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/IEventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/IEventStore.kt @@ -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() } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/ObservableEventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/ObservableEventStore.kt index d65002ebc4..35e0427f25 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/ObservableEventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/ObservableEventStore.kt @@ -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() } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventStore.kt index ce336b1e13..7e1c912d55 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventStore.kt @@ -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() } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/FullTextSearchModule.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/FullTextSearchModule.kt index 456e7d28c7..3de6859e4b 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/FullTextSearchModule.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/FullTextSearchModule.kt @@ -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( + 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 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 47e2ac3cfa..ad9088275b 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 @@ -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() } diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SearchTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SearchTest.kt index b79d0d2679..bdf7ae690a 100644 --- a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SearchTest.kt +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SearchTest.kt @@ -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 -> diff --git a/quartz/src/jvmMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/fs/FsEventStore.kt b/quartz/src/jvmMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/fs/FsEventStore.kt index e15f810a6a..0313eebda5 100644 --- a/quartz/src/jvmMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/fs/FsEventStore.kt +++ b/quartz/src/jvmMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/fs/FsEventStore.kt @@ -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//` 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 diff --git a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/fs/FsSearchTest.kt b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/fs/FsSearchTest.kt index 766644a93f..9251311df9 100644 --- a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/fs/FsSearchTest.kt +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/fs/FsSearchTest.kt @@ -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(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(Filter(search = "uniqnote")).map { it.id }) + assertEquals(listOf(long.id), store.query(Filter(search = "uniqlong")).map { it.id }) + } + @Test fun `reindexFullTextSearch ignores non-searchable kinds`() = runBlocking {