From c1d1f16fc088b6eb2a293b003470ff86186e00ee Mon Sep 17 00:00:00 2001 From: Bogdan Abaev Date: Fri, 4 Sep 2026 15:21:38 -0700 Subject: [PATCH] count progress on per-chunk basis Instead of per-item. It gives a more accurante representation of the progress since chunks are units of work and items will have varying amount of text to embed. Updated the advanced prefs pane with loading bars instead of per-library breakdown to reflect that. --- .../preferences/preferences_advanced.js | 52 +- .../preferences/preferences_advanced.xhtml | 13 +- chrome/content/zotero/xpcom/embeddings.js | 499 ++++++++++-------- chrome/locale/en-US/zotero/preferences.ftl | 3 + scss/preferences/_advanced.scss | 21 + test/tests/embeddingsTest.js | 128 ++++- 6 files changed, 451 insertions(+), 265 deletions(-) diff --git a/chrome/content/zotero/preferences/preferences_advanced.js b/chrome/content/zotero/preferences/preferences_advanced.js index 7bec34eb50..dc937268e2 100644 --- a/chrome/content/zotero/preferences/preferences_advanced.js +++ b/chrome/content/zotero/preferences/preferences_advanced.js @@ -163,9 +163,8 @@ Zotero_Preferences.Advanced = { // Phase / status message let phaseLabel = document.getElementById('semantic-search-phase'); - let hasRemaining = status.libraries.some( - lib => lib.indexed < lib.eligible - || lib.indexedAttachments < lib.eligibleAttachments); + let hasRemaining = status.items.done < status.items.total + || status.chunks.done < status.chunks.total; if (status.error) { document.l10n.setAttributes(phaseLabel, 'preferences-advanced-semantic-search-error', { error: status.error }); } @@ -225,32 +224,27 @@ Zotero_Preferences.Advanced = { stopButton.hidden = !status.indexing; stopButton.disabled = !!status.stopping; - // Per-library "indexed / total" counts - let grid = document.getElementById('semantic-search-libraries'); - if (grid.childElementCount !== status.libraries.length * 2) { - grid.textContent = ''; - for (let i = 0; i < status.libraries.length; i++) { - grid.append(document.createXULElement('label'), document.createXULElement('label')); - } - } - status.libraries.forEach((lib, i) => { - grid.children[i * 2].setAttribute( - 'value', - Zotero.Utilities.Internal.stringWithColon(lib.name) - ); - // Attachment fulltext is reported on its own, since it's a much - // larger and much slower job than the rest -- one combined count - // would look stalled. With fulltext indexing off, none are - // eligible and only the item count is shown. - let counts = `${lib.indexed.toLocaleString()} / ${lib.eligible.toLocaleString()}`; - if (lib.eligibleAttachments) { - counts += ` ${Zotero.getString('general.and')} ` - + `${lib.indexedAttachments.toLocaleString()} / ` - + `${lib.eligibleAttachments.toLocaleString()} ` - + Zotero.getString('itemTypes.attachment'); - } - grid.children[i * 2 + 1].setAttribute('value', counts); - }); + // One bar for items, notes and annotations, one for attachment full + // text + document.getElementById('semantic-search-items-row').hidden = !status.items.total; + this._updateSemanticSearchBar('items', status.items); + document.getElementById('semantic-search-attachments-row').hidden + = !Zotero.Prefs.get('embeddings.indexFulltext'); + this._updateSemanticSearchBar('attachments', status.chunks); + }, + + + // Fill one progress row with done out of total. The percentage is + // rounded down, so it reads 100% only when everything is done. + _updateSemanticSearchBar: function (name, { done, total }) { + let bar = document.getElementById(`semantic-search-${name}-progress`); + bar.max = Math.max(total, 1); + bar.value = done; + document.l10n.setAttributes( + document.getElementById(`semantic-search-${name}-value`), + 'preferences-advanced-semantic-search-progress-value', + { done, total, percent: total ? Math.floor(done / total * 100) : 0 } + ); }, diff --git a/chrome/content/zotero/preferences/preferences_advanced.xhtml b/chrome/content/zotero/preferences/preferences_advanced.xhtml index 17e1afcdfd..0eaa19b0de 100644 --- a/chrome/content/zotero/preferences/preferences_advanced.xhtml +++ b/chrome/content/zotero/preferences/preferences_advanced.xhtml @@ -343,7 +343,18 @@ data-l10n-id="preferences-advanced-semantic-search-stop" hidden="true"/> - + + + + + + diff --git a/chrome/content/zotero/xpcom/embeddings.js b/chrome/content/zotero/xpcom/embeddings.js index c92403cdb0..d21f6433b2 100644 --- a/chrome/content/zotero/xpcom/embeddings.js +++ b/chrome/content/zotero/xpcom/embeddings.js @@ -375,10 +375,9 @@ Zotero.Embeddings = new function () { await Zotero.DB.queryAsync("DROP TABLE IF EXISTS embeddings.itemChunkCounts"); // One row per chunk of an item's text, each carrying the hash of the // item's full source. For attachments, the block range and offsets - // locate the chunk's text in the document (see getMatchingChunks()), - // and one with no text gets a single NULL-embedding row so it counts - // as processed. No foreign key: deletions are handled by the notifier - // and eligibility pruning. + // locate the chunk's text in the document (see getMatchingChunks()). + // No foreign key: deletions are handled by the notifier and + // eligibility pruning. await Zotero.DB.queryAsync( "CREATE TABLE embeddings.itemEmbeddings (\n" + " itemID INTEGER NOT NULL,\n" @@ -414,13 +413,15 @@ Zotero.Embeddings = new function () { + ")" ); // How many chunks each attachment's current source splits into, - // recorded at extraction, so fulltext work is known without reading - // a pack. Other item types count as one apiece. + // recorded at extraction, and how many of them have rows stored. An + // attachment is indexed when the two are equal. Other item types + // count as one apiece. await Zotero.DB.queryAsync( "CREATE TABLE embeddings.itemChunkCounts (\n" + " itemID INTEGER PRIMARY KEY,\n" + " sourceHash TEXT NOT NULL,\n" - + " chunks INTEGER NOT NULL\n" + + " chunks INTEGER NOT NULL,\n" + + " embedded INTEGER NOT NULL DEFAULT 0\n" + ")" ); await Zotero.DB.queryAsync( @@ -1627,9 +1628,13 @@ Zotero.Embeddings.Indexing = new function () { // Attachments prepared so far ({ done, total }), while _phase is 'extracting' let _extractionProgress = null; let _lastError = null; - let _status = new Map(); // libraryID -> { name, indexed, eligible } + // Indexed and eligible items, notes and annotations, and fulltext work + // in chunks (see _getChunkCounts()) + let _itemCounts = { done: 0, total: 0 }; + let _chunkCounts = { done: 0, total: 0 }; let _progressListeners = new Set(); - let _lastStatusRefresh = 0; + let _lastTick = 0; + let _lastCountRefresh = 0; // The indexing queues. Producers (the item notifier, startIndexing()) only // add itemIDs here; _run() is the single consumer that drains them. @@ -1661,17 +1666,21 @@ Zotero.Embeddings.Indexing = new function () { // Wait longer than the usual debounce before retrying a run that was held // off for memory const LOW_MEMORY_RETRY_DELAY = 5 * 60 * 1000; - // Most often the per-library status counts are recomputed during a run, - // since each refresh is a pass over the database - const STATUS_REFRESH_INTERVAL = 5000; + // How often a run reports progress, and how often the per-library + // counts -- a scan of the index -- are recomputed within that + const PROGRESS_EMIT_INTERVAL = 1000; + const COUNT_REFRESH_INTERVAL = 5000; // Debounce before starting the consumer, so a burst of changes (e.g. an // import) is picked up in one pass const KICK_DELAY = 3000; - // itemIDs taken off a queue per pass of _indexItems(). Small, so the - // head of the queue is searchable soon: an item's rows are written only - // once all of its chunks are embedded, and a pass sorts every chunk in - // the slice by length before batching. + // itemIDs taken off a queue per pass of _indexItems(). A pass holds the + // text of every chunk in the slice and sorts them by length, so this + // bounds memory and how long the head of the queue waits behind the + // rest; the regular queue is checked between passes. const QUEUE_SLICE_SIZE = 32; + // Bump when chunking changes, so stored attachment rows are rebuilt (see + // _getAttachmentSourceHash()) + const CHUNKER_VERSION = 1; // The inference process's memory arena only grows: fragmentation from // varying batch shapes accumulates and is never returned to the OS @@ -1861,7 +1870,7 @@ Zotero.Embeddings.Indexing = new function () { // model except the newly-selected one. await _clearEmbeddings(); await Zotero.Embeddings.pruneModels(); - _status.clear(); + _clearCounts(); if (Zotero.Embeddings.isEnabled()) { Zotero.Embeddings.Indexing.startIndexing(); @@ -2242,21 +2251,22 @@ Zotero.Embeddings.Indexing = new function () { } } - // Items in a library that have a stored embedding -- the numerators for - // indexing progress, split the way _getEligibleItemIDs() splits the - // denominators. An item's chunks count as one item, and an attachment - // recorded as processed-but-empty counts as done (see _indexItems()). - async function _getIndexedCounts(libraryID) { - let attachmentTypeID = Zotero.ItemTypes.getID('attachment'); - let row = await Zotero.DB.rowQueryAsync( - "SELECT " - + "COUNT(DISTINCT CASE WHEN itemTypeID!=? THEN itemID END) AS items, " - + "COUNT(DISTINCT CASE WHEN itemTypeID=? THEN itemID END) AS attachments " - + "FROM embeddings.itemEmbeddings JOIN items USING (itemID) " - + "WHERE libraryID=?", - [attachmentTypeID, attachmentTypeID, libraryID] + // Indexed items, notes and annotations -- the numerator for their + // progress; an item's chunks count as one item. Attachments are measured + // in chunks instead (see _getChunkCounts()). + async function _getIndexedItemCount() { + return Zotero.DB.valueQueryAsync( + "SELECT COUNT(DISTINCT itemID) FROM embeddings.itemEmbeddings " + + "JOIN items USING (itemID) WHERE itemTypeID!=?", + Zotero.ItemTypes.getID('attachment') ); - return { items: row.items, attachments: row.attachments }; + } + + // Counts after the stored index is cleared: nothing done, and no chunk + // counts until attachments are extracted again + function _clearCounts() { + _itemCounts = { done: 0, total: _itemCounts.total }; + _chunkCounts = { done: 0, total: 0 }; } @@ -2350,13 +2360,11 @@ Zotero.Embeddings.Indexing = new function () { // The staleness key for an attachment's stored chunks, standing in for // the text hash other item types use. Derived from the file's identity // (path, size, mtime) rather than its extracted text, so the skip check - // every indexing pass runs costs a stat rather than an extraction. A - // change to the extraction or chunking logic isn't detected -- rebuild - // the index after one. (Nor is a pack regenerated for a processor bump - // without the file changing -- the vectors stay derived from the older - // extraction until the file changes or the index is rebuilt.) - // Returns null when the attachment has no readable file, which also - // means there's nothing to extract. + // every indexing pass runs costs a stat rather than an extraction. The + // extractor's version is part of it, and so is CHUNKER_VERSION: stored + // rows are resumed by chunk index, so a change to how text is chunked + // has to invalidate them. Returns null when the attachment has no + // readable file, which also means there's nothing to extract. async function _getAttachmentSourceHash(item) { try { let path = await item.getFilePathAsync(); @@ -2364,13 +2372,9 @@ Zotero.Embeddings.Indexing = new function () { return null; } let { size, lastModified } = await IOUtils.stat(path); - // The extractor's identity is part of the source: a processor - // upgrade changes what the same file extracts to, and the stored - // rows point into that extraction (blocks, offsets), so they go - // stale with it just as with a changed file let extractor = await Zotero.SDT.getProcessorVersion(item); return Zotero.Utilities.Internal.md5( - [path, size, lastModified, extractor].join('|')); + [path, size, lastModified, extractor, CHUNKER_VERSION].join('|')); } catch (e) { if (e.name !== 'NotFoundError') { @@ -2479,18 +2483,18 @@ Zotero.Embeddings.Indexing = new function () { })); } - // The stored source hash of each of the given items that has rows in - // the given table, read in one query per chunk rather than one per item - // (every start re-enqueues the whole library to find what changed). - // Every chunk row of an item carries the same hash. - async function _getStoredHashes(itemIDs, table = 'itemEmbeddings') { + // The stored source hash of each of the given items that has rows, read + // in one query per chunk rather than one per item (every start + // re-enqueues the whole library to find what changed). Every chunk row + // of an item carries the same hash. + async function _getStoredHashes(itemIDs) { let storedHashes = new Map(); let chunkSize = 500; for (let i = 0; i < itemIDs.length; i += chunkSize) { let chunk = itemIDs.slice(i, i + chunkSize); let rows = await Zotero.DB.queryAsync( - "SELECT DISTINCT itemID, sourceHash FROM embeddings." + table - + " WHERE itemID IN (" + chunk.map(() => '?').join(',') + ")", + "SELECT DISTINCT itemID, sourceHash FROM embeddings.itemEmbeddings " + + "WHERE itemID IN (" + chunk.map(() => '?').join(',') + ")", chunk ); for (let row of rows) { @@ -2500,36 +2504,100 @@ Zotero.Embeddings.Indexing = new function () { return storedHashes; } - // How long between extraction-progress emissions - const EXTRACTION_EMIT_INTERVAL = 1000; - - // The given attachments whose stored embeddings are stale or absent -- - // the ones the embedding pass will read packs for. An up-to-date + // The given attachments whose stored embeddings are incomplete or stale + // -- the ones the embedding pass will read packs for. An indexed // attachment won't be re-embedded and a fileless one has nothing to // extract, so neither needs a pack. async function _staleAttachments(itemIDs, shouldStop) { let items = await Zotero.Items.getAsync(itemIDs); - let storedHashes = await _getStoredHashes(itemIDs); + let ledger = await _getChunkCountRows(itemIDs); let stale = []; for (let item of items) { if (shouldStop()) { break; } let hash = await _getAttachmentSourceHash(item); - if (hash && storedHashes.get(item.id) !== hash) { + if (hash && !_isIndexed(ledger.get(item.id), hash)) { stale.push({ item, hash }); } } return stale; } - // Record how many chunks an attachment's current source splits into (see - // the itemChunkCounts table in _setUpDB()) - async function _storeChunkCount(itemID, hash, chunks) { + // The ledger rows (see itemChunkCounts in _setUpDB()) of the given + // attachments, as itemID -> { sourceHash, chunks, embedded } + async function _getChunkCountRows(itemIDs) { + let ledger = new Map(); + let chunkSize = 500; + for (let i = 0; i < itemIDs.length; i += chunkSize) { + let chunk = itemIDs.slice(i, i + chunkSize); + let rows = await Zotero.DB.queryAsync( + "SELECT itemID, sourceHash, chunks, embedded FROM embeddings.itemChunkCounts " + + "WHERE itemID IN (" + chunk.map(() => '?').join(',') + ")", + chunk + ); + for (let row of rows) { + ledger.set(row.itemID, row); + } + } + return ledger; + } + + // Whether a ledger row says every chunk of the source `hash` is stored + function _isIndexed(row, hash) { + return !!row && row.sourceHash === hash && row.embedded >= row.chunks; + } + + // Record how many chunks an attachment's current source splits into and + // how many are stored + async function _storeChunkCount(itemID, hash, chunks, embedded = 0) { await Zotero.DB.queryAsync( - "REPLACE INTO embeddings.itemChunkCounts (itemID, sourceHash, chunks) " - + "VALUES (?, ?, ?)", - [itemID, hash, chunks] + "REPLACE INTO embeddings.itemChunkCounts (itemID, sourceHash, chunks, embedded) " + + "VALUES (?, ?, ?, ?)", + [itemID, hash, chunks, embedded] + ); + } + + // Recount an attachment's stored chunks in the ledger, from the rows + async function _updateEmbeddedCount(itemID, hash) { + await Zotero.DB.queryAsync( + "UPDATE embeddings.itemChunkCounts SET embedded=(" + + "SELECT COUNT(*) FROM embeddings.itemEmbeddings " + + "WHERE itemID=? AND sourceHash=? AND embedding IS NOT NULL" + + ") WHERE itemID=?", + [itemID, hash, itemID] + ); + } + + // Store one chunk's vector. Source references are stored only for + // attachments, whose text lives in a file: they're what the search + // preview is re-derived from (see getMatchingChunks()). Other item types + // are their own preview. + async function _insertChunkRow(entry, chunkIndex, vector) { + let chunk = entry.chunks[chunkIndex]; + let isAttachment = entry.item.isAttachment(); + let blob = new Uint8Array(vector.buffer, vector.byteOffset, vector.byteLength); + // Keep the embedding blob out of debug output + await Zotero.DB.queryAsync( + "INSERT INTO embeddings.itemEmbeddings " + + "(itemID, chunkIndex, embedding, sourceHash, " + + "startBlock, endBlock, startOffset, endOffset, " + + "textCheck, sectionPart, sectionParts) " + + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", + [ + entry.item.id, + chunkIndex, + blob, + entry.hash, + chunk.startBlock ?? null, + chunk.endBlock ?? null, + isAttachment ? chunk.startOffset ?? null : null, + isAttachment ? chunk.endOffset ?? null : null, + isAttachment ? Zotero.Embeddings.textCheck(chunk.text) : null, + chunk.sectionPart ?? null, + chunk.sectionParts ?? null + ], + { debugParams: false } ); } @@ -2591,48 +2659,37 @@ Zotero.Embeddings.Indexing = new function () { } // Extract the structured text of the attachments about to be embedded, so - // the embedding step reads cached packs instead of extracting inline. - // `extracted` collects what this run has put through extraction -- - // attempts included, so a failure isn't retried for the rest of the run - // -- and lets an interrupted pass resume without re-reading packs. + // the embedding step reads cached packs instead of extracting inline. An + // attachment whose chunks are counted for its current source was + // extracted then and is skipped, so a resumed run has nothing to prepare. // // Runs with the engine shut down: extraction takes a long while, and the // model would otherwise sit in memory throughout, competing with the // worker for the same cores. - async function _extractAttachments(itemIDs, extracted, shouldStop) { - let toExtract = (await _staleAttachments(itemIDs, shouldStop)) - .filter(({ item }) => !extracted.has(item.id)); + async function _extractAttachments(itemIDs, shouldStop) { + let stale = await _staleAttachments(itemIDs, shouldStop); + let counted = await _getChunkCountRows(stale.map(({ item }) => item.id)); + let toExtract = stale.filter(({ item, hash }) => counted.get(item.id)?.sourceHash !== hash); if (!toExtract.length) { return; } _setPhase('extracting'); await Zotero.Embeddings.shutdownEngine({ modelChanged: false }); - // Chunk counts are taken while the pack is fresh, through the same - // derivation the embedder uses so the two can't disagree. One already - // recorded for the current source is kept -- chunking again would - // only repeat it. - let counted = await _getStoredHashes( - toExtract.map(({ item }) => item.id), 'itemChunkCounts'); let progress = { done: 0, total: toExtract.length }; _extractionProgress = progress; _emitProgress(); - let lastEmit = Date.now(); try { for (let { item, hash } of toExtract) { if (shouldStop()) { return; } - extracted.add(item.id); await Zotero.SDT.ensure(item.id); - if (counted.get(item.id) !== hash) { - let chunks = await _getAttachmentChunks(item); - await _storeChunkCount(item.id, hash, chunks ? chunks.length : 0); - } + // Counted while the pack is fresh, through the same derivation + // the embedder uses so the two can't disagree + let chunks = await _getAttachmentChunks(item); + await _storeChunkCount(item.id, hash, chunks ? chunks.length : 0); progress.done++; - if (Date.now() - lastEmit >= EXTRACTION_EMIT_INTERVAL) { - lastEmit = Date.now(); - _emitProgress(); - } + await _tick(); } } finally { @@ -2648,7 +2705,8 @@ Zotero.Embeddings.Indexing = new function () { // // @param {Zotero.Item[]} items // @param {Object} [options] - // @param {Function} [options.onProgress] - Called as { done, total } + // @param {Function} [options.onProgress] - Called after every batch as + // { done, total } items // @param {Number} [options.maxBatchItems=20] - Most items per engine call // @param {Number} [options.batchTokenBudget=3000] - Most tokens per engine // call, counting every text in the batch as long as its longest one, @@ -2665,6 +2723,8 @@ Zotero.Embeddings.Indexing = new function () { await Zotero.Items.loadDataTypes(items, ['itemData', 'note', 'annotation']); let storedHashes = await _getStoredHashes(items.map(item => item.id)); + let ledger = await _getChunkCountRows( + items.filter(item => item.isAttachment()).map(item => item.id)); let toEmbed = []; let toDelete = []; @@ -2672,16 +2732,17 @@ Zotero.Embeddings.Indexing = new function () { // An attachment's text lives in a file, so its staleness check is // a file-identity hash rather than a text hash -- reading and // sectioning every attachment on every pass would defeat the - // check's purpose + // check's purpose -- and the ledger says whether all of its + // chunks are stored if (item.isAttachment()) { let hash = await _getAttachmentSourceHash(item); if (!hash) { - if (storedHashes.has(item.id)) { + if (ledger.has(item.id) || storedHashes.has(item.id)) { toDelete.push(item.id); } continue; } - if (storedHashes.get(item.id) !== hash) { + if (!_isIndexed(ledger.get(item.id), hash)) { toEmbed.push({ item, hash }); } continue; @@ -2708,7 +2769,6 @@ Zotero.Embeddings.Indexing = new function () { // abstract, or an annotation's passage and comment, fit the window in // almost all cases, so they're embedded as a single chunk and the // pipeline truncates the rare outlier. - let emptyAttachments = []; for (let entry of toEmbed) { // Extracting an attachment and tokenizing a chunk's worth of long // notes take real time, and a stop request can be a model switch @@ -2722,13 +2782,11 @@ Zotero.Embeddings.Indexing = new function () { if (entry.item.isAttachment()) { entry.chunks = await _getAttachmentChunks(entry.item); // Nothing embeddable anywhere in the attachment (missing - // file, password-protected, no text layer). Record the - // attempt anyway, so the item counts as processed and isn't - // looked at again until the file changes. - if (!entry.chunks || !entry.chunks.length) { + // file, password-protected, no text layer). Recorded below + // anyway, so the item counts as processed and isn't looked + // at again until the file changes. + if (!entry.chunks) { entry.chunks = []; - emptyAttachments.push(entry); - continue; } } else if (entry.item.isNote()) { @@ -2745,47 +2803,52 @@ Zotero.Embeddings.Indexing = new function () { } entry.vectors = new Array(entry.chunks.length); entry.remaining = entry.chunks.length; + entry.stored = new Set(); } - // An attachment with nothing to embed is still processed: replace - // whatever an older file left with a single embedding-less row - // carrying the current source hash, so the indexed count converges on - // the eligible count instead of these items reading as forever - // unindexed (see _setUpDB()) - if (emptyAttachments.length) { + // Prepare the attachments' rows: drop what an older source or a + // longer text left, keep the current source's rows so an interrupted + // item resumes where it stopped, and bring the ledger in line. One + // with nothing to embed is complete at zero chunks. + let attachmentEntries = toEmbed.filter(entry => entry.item.isAttachment()); + if (attachmentEntries.length) { await Zotero.DB.executeTransaction(async function () { - for (let entry of emptyAttachments) { + for (let entry of attachmentEntries) { // The item may have been deleted while we were extracting if (!Zotero.Items.get(entry.item.id)) { continue; } await Zotero.DB.queryAsync( - "DELETE FROM embeddings.itemEmbeddings WHERE itemID=?", + "DELETE FROM embeddings.itemEmbeddings " + + "WHERE itemID=? AND (sourceHash!=? OR chunkIndex>=?)", + [entry.item.id, entry.hash, entry.chunks.length] + ); + entry.stored = new Set(await Zotero.DB.columnQueryAsync( + "SELECT chunkIndex FROM embeddings.itemEmbeddings WHERE itemID=?", entry.item.id - ); - await Zotero.DB.queryAsync( - "INSERT INTO embeddings.itemEmbeddings " - + "(itemID, chunkIndex, embedding, sourceHash) " - + "VALUES (?, 0, NULL, ?)", - [entry.item.id, entry.hash] - ); - await _storeChunkCount(entry.item.id, entry.hash, 0); + )); + entry.remaining -= entry.stored.size; + await _storeChunkCount( + entry.item.id, entry.hash, entry.chunks.length, entry.stored.size); } }); } - toEmbed = toEmbed.filter(entry => entry.chunks.length); + toEmbed = toEmbed.filter(entry => entry.remaining > 0); // What gets embedded is the chunk's embedText (its text plus any // outline-path context); plain chunks embed their text as is. let embedText = chunk => chunk.embedText || chunk.text; // Batches are packed from the flattened chunks, so an item's chunks - // can span batches; its rows are written only once every chunk's - // vector is in, as one transaction, so a stop mid-item never leaves a - // partial set that the source hash would report as complete + // can span batches. An attachment's rows are written as each batch + // finishes, since the ledger says when it's complete; any other + // item's are written once every chunk is in, since its rows' hash + // alone marks it done. let units = []; for (let entry of toEmbed) { for (let chunkIndex = 0; chunkIndex < entry.chunks.length; chunkIndex++) { - units.push({ entry, chunkIndex }); + if (!entry.stored.has(chunkIndex)) { + units.push({ entry, chunkIndex }); + } } } // Pack batches from chunks of similar size, measured in the model's @@ -2837,81 +2900,50 @@ Zotero.Embeddings.Indexing = new function () { Zotero.debug('Embeddings: engine restarted to apply new thread count'); } let completed = []; - for (let j = 0; j < batch.length; j++) { - let { entry, chunkIndex } = batch[j]; - entry.vectors[chunkIndex] = vectors[j]; - if (--entry.remaining === 0) { - completed.push(entry); - } - } - if (completed.length) { - await Zotero.DB.executeTransaction(async function () { - for (let entry of completed) { + await Zotero.DB.executeTransaction(async function () { + let touched = new Set(); + for (let j = 0; j < batch.length; j++) { + let { entry, chunkIndex } = batch[j]; + if (entry.item.isAttachment()) { // The item may have been deleted while the batch was // embedding -- don't write its vectors back after the // delete notifier removed them - if (!Zotero.Items.get(entry.item.id)) { - continue; - } - // Replace the item's rows as a unit, so a previously - // longer text never leaves stale chunks behind - await Zotero.DB.queryAsync( - "DELETE FROM embeddings.itemEmbeddings WHERE itemID=?", - entry.item.id - ); - for (let k = 0; k < entry.vectors.length; k++) { - let vector = entry.vectors[k]; - let chunk = entry.chunks[k]; - let blob = new Uint8Array( - vector.buffer, vector.byteOffset, vector.byteLength - ); - // Source references are stored only for - // attachments, whose text lives in a file: they're - // what the search-results preview is re-derived - // from, with textCheck fingerprinting the text - // they pointed to when embedded (see - // getMatchingChunks()). Other item types are - // their own preview. - // Keep the embedding blobs out of debug output. - let isAttachment = entry.item.isAttachment(); - await Zotero.DB.queryAsync( - "INSERT INTO embeddings.itemEmbeddings " - + "(itemID, chunkIndex, embedding, sourceHash, " - + "startBlock, endBlock, startOffset, endOffset, " - + "textCheck, sectionPart, sectionParts) " - + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", - [ - entry.item.id, - k, - blob, - entry.hash, - chunk.startBlock ?? null, - chunk.endBlock ?? null, - isAttachment ? chunk.startOffset ?? null : null, - isAttachment ? chunk.endOffset ?? null : null, - isAttachment - ? Zotero.Embeddings.textCheck(chunk.text) - : null, - chunk.sectionPart ?? null, - chunk.sectionParts ?? null - ], - { debugParams: false } - ); - } - // Normally recorded at extraction; written here too - // for an attachment that reached embedding without - // passing through it, so the count always matches - // the rows - if (entry.item.isAttachment()) { - await _storeChunkCount(entry.item.id, entry.hash, entry.chunks.length); + if (Zotero.Items.get(entry.item.id)) { + await _insertChunkRow(entry, chunkIndex, vectors[j]); + touched.add(entry); } } - }); + else { + entry.vectors[chunkIndex] = vectors[j]; + } + if (--entry.remaining === 0) { + completed.push(entry); + } + } + for (let entry of touched) { + await _updateEmbeddedCount(entry.item.id, entry.hash); + } + for (let entry of completed) { + if (entry.item.isAttachment() || !Zotero.Items.get(entry.item.id)) { + continue; + } + // Replace the item's rows as a unit, so a previously + // longer text never leaves stale chunks behind + await Zotero.DB.queryAsync( + "DELETE FROM embeddings.itemEmbeddings WHERE itemID=?", + entry.item.id + ); + for (let k = 0; k < entry.vectors.length; k++) { + await _insertChunkRow(entry, k, entry.vectors[k]); + } + } + }); + if (completed.length) { _notifyIndexed(completed.map(entry => entry.item.id)); done += completed.length; - if (onProgress) { - onProgress({ done, total: toEmbed.length }); - } + } + if (onProgress) { + onProgress({ done, total: toEmbed.length }); } // Yield so the UI thread stays responsive between batches await Zotero.Promise.delay(0); @@ -2949,10 +2981,9 @@ Zotero.Embeddings.Indexing = new function () { /** * Current runner state, for the preferences UI. * - * Each library's counts come in two disjoint pairs: `indexed`/`eligible` - * for items, notes, and annotations, and `indexedAttachments`/ - * `eligibleAttachments` for attachment fulltext, which is a far bigger and - * slower job. Callers that want whole-library coverage add them up. + * Progress comes in two disjoint pairs: `items` counts items, notes and + * annotations, and `chunks` measures attachment fulltext -- a far bigger + * and slower job -- in chunks stored so far. */ this.getStatus = function () { return { @@ -2968,8 +2999,9 @@ Zotero.Embeddings.Indexing = new function () { queued: { items: _queue.size, attachments: _attachmentQueue.size }, downloadProgress: _downloadProgress, extractionProgress: _extractionProgress, - error: _lastError ? (_lastError.message || String(_lastError)) : null, - libraries: [..._status.entries()].map(([libraryID, s]) => ({ libraryID, ...s })) + items: _itemCounts, + chunks: _chunkCounts, + error: _lastError ? (_lastError.message || String(_lastError)) : null }; }; @@ -3003,8 +3035,9 @@ Zotero.Embeddings.Indexing = new function () { } /** - * Recompute per-library indexed/eligible counts (without indexing anything) - * and notify listeners. Used by the preferences UI to show current state. + * Recompute every count and notify listeners. The eligibility pass makes + * this the expensive refresh, for the pane opening and a run's ends; a + * run in progress ticks with _tick() instead. * * @return {Promise} - The status object */ @@ -3017,25 +3050,49 @@ Zotero.Embeddings.Indexing = new function () { } await Zotero.Embeddings.initDB(); let eligibleByLibrary = await _getEligibleItemIDs(); + let total = 0; for (let library of _indexableLibraries()) { - let eligible = eligibleByLibrary.get(library.libraryID) - || { items: [], attachments: [] }; - let indexed = await _getIndexedCounts(library.libraryID); - // Attachments are counted separately from everything else, not - // included in it -- the two pairs are disjoint, and consumers sum - // them when they want the whole library - _status.set(library.libraryID, { - name: library.name, - indexed: indexed.items, - eligible: eligible.items.length, - indexedAttachments: indexed.attachments, - eligibleAttachments: eligible.attachments.length - }); + total += eligibleByLibrary.get(library.libraryID)?.items.length || 0; } + _itemCounts = { done: await _getIndexedItemCount(), total }; + _chunkCounts = await _getChunkCounts(); + _lastCountRefresh = Date.now(); _emitProgress(); return Zotero.Embeddings.Indexing.getStatus(); }; + // Fulltext work in chunks, from the ledger + async function _getChunkCounts() { + let row = await Zotero.DB.rowQueryAsync( + "SELECT COALESCE(SUM(chunks), 0) AS total, COALESCE(SUM(embedded), 0) AS done " + + "FROM embeddings.itemChunkCounts" + ); + return { done: row.done, total: row.total }; + } + + // Progress tick for a run's inner loops: refresh the chunk counts and + // emit, at most once per PROGRESS_EMIT_INTERVAL. The indexed item count + // costs a scan of the index, so it's recomputed less often; the eligible + // count comes from the last full refresh. + async function _tick() { + let now = Date.now(); + if (now - _lastTick < PROGRESS_EMIT_INTERVAL) { + return; + } + _lastTick = now; + try { + _chunkCounts = await _getChunkCounts(); + if (now - _lastCountRefresh >= COUNT_REFRESH_INTERVAL) { + _lastCountRefresh = now; + _itemCounts = { done: await _getIndexedItemCount(), total: _itemCounts.total }; + } + } + catch (e) { + Zotero.logError(e); + } + _emitProgress(); + } + /** * Start (or resume) indexing: clear a previous stopIndexing(), drop stored * embeddings for items that no longer have indexable text, re-enqueue @@ -3093,7 +3150,7 @@ Zotero.Embeddings.Indexing = new function () { Zotero.debug(`Embeddings: stored embeddings are from '${indexed || 'unknown'}' ` + `but the active model is '${current}' -- clearing for reindexing`); await _clearEmbeddings(); - _status.clear(); + _clearCounts(); } await Zotero.DB.queryAsync( "REPLACE INTO embeddings.itemEmbeddingsMeta (key, value) VALUES ('modelVersion', ?)", @@ -3140,6 +3197,7 @@ Zotero.Embeddings.Indexing = new function () { _indexing = true; _stopping = false; _lastError = null; + _lastTick = 0; _startProcMonitor(); _startIdleWatch(); try { @@ -3163,11 +3221,8 @@ Zotero.Embeddings.Indexing = new function () { let indexOptions = () => ({ shouldStop, batchTokenBudget: _tokenBudget, - onProgress: () => _refreshStatusThrottled() + onProgress: () => _tick() }); - // Attachments this run has already extracted, so a cycle the - // regular queue interrupted resumes without re-reading packs - let extracted = new Set(); // One step at a time, in a fixed order: the regular queue, then // the attachment queue's extraction, then its embedding. The // regular queue goes first and preempts the attachments, so a @@ -3181,7 +3236,7 @@ Zotero.Embeddings.Indexing = new function () { } let itemIDs = [..._attachmentQueue]; _attachmentQueue.clear(); - await _extractAttachments(itemIDs, extracted, shouldStop); + await _extractAttachments(itemIDs, shouldStop); await _embedAttachments(itemIDs, shouldStop, indexOptions); } await Zotero.Embeddings.Indexing.refreshStatus(); @@ -3197,6 +3252,15 @@ Zotero.Embeddings.Indexing = new function () { _phase = 'idle'; _downloadProgress = null; _extractionProgress = null; + // A run cut short by a stop or an error still reports what's + // stored + try { + _chunkCounts = await _getChunkCounts(); + _itemCounts = { done: await _getIndexedItemCount(), total: _itemCounts.total }; + } + catch (e) { + Zotero.logError(e); + } _emitProgress(); // Pick up anything enqueued while we were finishing up if ((_queue.size || _attachmentQueue.size) && !_stopping) { @@ -3216,15 +3280,6 @@ Zotero.Embeddings.Indexing = new function () { } } - function _refreshStatusThrottled() { - let now = Date.now(); - if (now - _lastStatusRefresh < STATUS_REFRESH_INTERVAL) { - return; - } - _lastStatusRefresh = now; - Zotero.Embeddings.Indexing.refreshStatus().catch(e => Zotero.logError(e)); - } - this.stopIndexing = function () { _stopping = true; _queue.clear(); diff --git a/chrome/locale/en-US/zotero/preferences.ftl b/chrome/locale/en-US/zotero/preferences.ftl index 3f6f2f95b2..f2ff117702 100644 --- a/chrome/locale/en-US/zotero/preferences.ftl +++ b/chrome/locale/en-US/zotero/preferences.ftl @@ -126,6 +126,9 @@ preferences-advanced-semantic-search-resume = .label = Resume preferences-advanced-semantic-search-stop = .label = Stop +preferences-advanced-semantic-search-items = Items and notes +preferences-advanced-semantic-search-attachments = Attachments +preferences-advanced-semantic-search-progress-value = { $percent }% ({ $done } / { $total }) preferences-advanced-semantic-search-switch-title = Change Mode? preferences-advanced-semantic-search-switch-text = Changing the mode will rebuild the search index, which can take a long time for large libraries. preferences-advanced-semantic-search-switch-button = Change Mode diff --git a/scss/preferences/_advanced.scss b/scss/preferences/_advanced.scss index 9ab813f512..b907aa8168 100644 --- a/scss/preferences/_advanced.scss +++ b/scss/preferences/_advanced.scss @@ -98,6 +98,27 @@ margin-inline-start: 6px; } +#semantic-search-progress { + display: grid; + grid-template-columns: max-content 1fr max-content; + align-items: center; + column-gap: 8px; + row-gap: 4px; + margin-top: 4px; + + .semantic-search-progress-row { + display: contents; + + &[hidden] { + display: none; + } + } + + progress { + width: 100%; + } +} + #db-maintenance-options { display: flex; gap: 6px; diff --git a/test/tests/embeddingsTest.js b/test/tests/embeddingsTest.js index e4648e5e1a..e68c698089 100644 --- a/test/tests/embeddingsTest.js +++ b/test/tests/embeddingsTest.js @@ -1632,6 +1632,18 @@ describe("Zotero.Embeddings", function () { sinon.stub(Zotero.SDT, 'ensure').resolves(true), sinon.stub(Zotero.SDT, 'getSections').resolves({ ok: true, sections }) ]; + // Status reports the ledger, and items apart from it + let assertStatusMatchesLedger = async () => { + let { items, chunks } = Zotero.Embeddings.Indexing.getStatus(); + assert.isAtLeast(items.done, 1); + assert.isAtLeast(items.total, items.done); + let ledger = await Zotero.DB.rowQueryAsync( + "SELECT COALESCE(SUM(chunks), 0) AS total, COALESCE(SUM(embedded), 0) AS done " + + "FROM embeddings.itemChunkCounts" + ); + assert.deepEqual(chunks, { done: ledger.done, total: ledger.total }); + return chunks; + }; try { Zotero.Prefs.set('embeddings.indexFulltext', true); await Zotero.Embeddings.Indexing.startIndexing(); @@ -1640,6 +1652,8 @@ describe("Zotero.Embeddings", function () { attachment.id ); assert.isAbove(counted, 1); + let chunks = await assertStatusMatchesLedger(); + assert.isAtLeast(chunks.total - chunks.done, counted); assert.equal(await Zotero.DB.valueQueryAsync( "SELECT COUNT(*) FROM embeddings.itemEmbeddings WHERE itemID=?", attachment.id @@ -1650,8 +1664,21 @@ describe("Zotero.Embeddings", function () { item.id ), 0); + // The count marks the attachment as extracted, so the next + // run goes straight to embedding it failAttachment = false; - await Zotero.Embeddings.Indexing.startIndexing(); + Zotero.SDT.ensure.resetHistory(); + let phases = new Set(); + let listener = status => phases.add(status.phase); + Zotero.Embeddings.Indexing.addProgressListener(listener); + try { + await Zotero.Embeddings.Indexing.startIndexing(); + } + finally { + Zotero.Embeddings.Indexing.removeProgressListener(listener); + } + assert.isFalse(Zotero.SDT.ensure.calledWith(attachment.id)); + assert.notInclude([...phases], 'extracting'); assert.equal(await Zotero.DB.valueQueryAsync( "SELECT chunks FROM embeddings.itemChunkCounts WHERE itemID=?", attachment.id @@ -1661,6 +1688,83 @@ describe("Zotero.Embeddings", function () { + "WHERE itemID=? AND embedding IS NOT NULL", attachment.id ), counted); + await assertStatusMatchesLedger(); + } + finally { + stubs.forEach(stub => stub.restore()); + Zotero.Prefs.clear('embeddings.indexFulltext'); + } + }); + + it("should store an attachment's chunks as they're embedded and resume from them", async function () { + this.timeout(60000); + let item = await createDataObject('item', { title: 'Parent of resumed attachment' }); + let attachment = await importPDFAttachment(item); + + let vector = new Float32Array(4).fill(0.5); + // More chunks than fit one engine call, so the attachment spans + // several batches + let sections = []; + for (let i = 0; i < 12; i++) { + sections.push(sdtSection('', i, ['Owls hunt at night. '.repeat(200)])); + } + // The first batch with the attachment's text succeeds, the next + // one fails, so the run ends with the attachment partly stored + let hunts = 0; + let failSecond = true; + let embedded = []; + let stubs = [ + sinon.stub(Zotero.Embeddings, 'embedPassages').callsFake(async (texts) => { + if (texts.some(text => text.includes('hunt'))) { + if (failSecond && ++hunts === 2) { + throw new Error('Embedding failed'); + } + embedded.push(...texts.filter(text => text.includes('hunt'))); + } + return texts.map(() => vector); + }), + sinon.stub(Zotero.Embeddings, 'isEnabled').returns(true), + sinon.stub(Zotero.Embeddings, 'getModelVersion').returns('test-model/1'), + sinon.stub(Zotero.Embeddings, 'isDownloaded').resolves(true), + sinon.stub(Zotero.Embeddings, 'download').resolves(), + sinon.stub(Zotero.Embeddings, 'ensureCalibration').resolves(), + sinon.stub(Zotero.Embeddings, 'getModelName').returns('bge-small-en-v1.5'), + sinon.stub(Zotero.SDT, 'ensure').resolves(true), + sinon.stub(Zotero.SDT, 'getSections').resolves({ ok: true, sections }) + ]; + let ledgerRow = () => Zotero.DB.rowQueryAsync( + "SELECT chunks, embedded FROM embeddings.itemChunkCounts WHERE itemID=?", + attachment.id + ); + let storedRows = () => Zotero.DB.valueQueryAsync( + "SELECT COUNT(*) FROM embeddings.itemEmbeddings " + + "WHERE itemID=? AND embedding IS NOT NULL", + attachment.id + ); + try { + Zotero.Prefs.set('embeddings.indexFulltext', true); + await Zotero.Embeddings.Indexing.startIndexing(); + let { chunks, embedded: stored } = await ledgerRow(); + assert.isAbove(chunks, 20); + assert.isAbove(stored, 0); + assert.isBelow(stored, chunks); + assert.equal(await storedRows(), stored); + assert.equal(embedded.length, stored); + + // The next run embeds only what's missing + failSecond = false; + embedded = []; + await Zotero.Embeddings.Indexing.startIndexing(); + assert.equal(embedded.length, chunks - stored); + let after = await ledgerRow(); + assert.equal(after.chunks, chunks); + assert.equal(after.embedded, chunks); + assert.equal(await storedRows(), chunks); + + // ...and a complete attachment isn't read again + let reads = Zotero.SDT.getSections.callCount; + await Zotero.Embeddings.Indexing.startIndexing(); + assert.equal(Zotero.SDT.getSections.callCount, reads); } finally { stubs.forEach(stub => stub.restore()); @@ -2082,22 +2186,20 @@ describe("Zotero.Embeddings", function () { Zotero.Prefs.set('embeddings.indexFulltext', true); await Zotero.Embeddings.Indexing.startIndexing(); - // The attempt is recorded as a single embedding-less row, so - // the item counts as processed and the progress counts align - let rows = await Zotero.DB.queryAsync( - "SELECT embedding, sourceHash FROM embeddings.itemEmbeddings WHERE itemID=?", + // The attempt is recorded as a complete ledger row of zero + // chunks, so the item counts as processed with nothing stored + let row = await Zotero.DB.rowQueryAsync( + "SELECT chunks, embedded FROM embeddings.itemChunkCounts WHERE itemID=?", attachment.id ); - assert.lengthOf(rows, 1); - assert.isNull(rows[0].embedding); - assert.ok(rows[0].sourceHash); - // Sections are read once to count chunks at extraction and - // once to embed - assert.equal(ourCalls(), 2); + assert.equal(row.chunks, 0); + assert.equal(row.embedded, 0); assert.equal(await Zotero.DB.valueQueryAsync( - "SELECT chunks FROM embeddings.itemChunkCounts WHERE itemID=?", + "SELECT COUNT(*) FROM embeddings.itemEmbeddings WHERE itemID=?", attachment.id ), 0); + // Sections are read once, to count chunks at extraction + assert.equal(ourCalls(), 1); // A processed-but-empty item can't be scored, and doesn't // break scoring for anything else @@ -2107,7 +2209,7 @@ describe("Zotero.Embeddings", function () { // The record makes later passes skip the attachment without // re-extracting, until the file changes await Zotero.Embeddings.Indexing.startIndexing(); - assert.equal(ourCalls(), 2); + assert.equal(ourCalls(), 1); } finally { stubs.forEach(stub => stub.restore());