From 69dac08246c3d1026661871b5cb20f9dd3d519ad Mon Sep 17 00:00:00 2001 From: uommou Date: Mon, 5 Oct 2026 13:38:50 +0900 Subject: [PATCH] [ZEPPELIN-6412] Shard embedding search index persistence by note Persist embedding search entries in per-note shard files so that updating a paragraph rewrites only the affected note instead of the entire index. Add note-level dirty tracking and locking, atomic shard writes, targeted recovery for missing or corrupt shards, orphan cleanup, and migration from the legacy single-file index. Add tests covering shard isolation, migration, recovery, deletion races, and partial-load failures. --- docs/embedding-search.md | 26 +- .../zeppelin/search/EmbeddingSearch.java | 672 +++++++++++---- .../search/EmbeddingSearchShardingTest.java | 786 ++++++++++++++++++ 3 files changed, 1333 insertions(+), 151 deletions(-) create mode 100644 zeppelin-server/src/test/java/org/apache/zeppelin/search/EmbeddingSearchShardingTest.java diff --git a/docs/embedding-search.md b/docs/embedding-search.md index 90fb266bb3e..2609687fd2d 100644 --- a/docs/embedding-search.md +++ b/docs/embedding-search.md @@ -59,8 +59,8 @@ finding the right query becomes a significant productivity bottleneck. │ 1. Embed query → cosine sim → find tables │ │ │ 2. Re-rank with table boost → top-20 │ │ │ ▼ │ -│ Index: text + title + output + tables embedding_index.bin│ -│ (persisted to disk, versioned) │ +│ Index: text + title + output + tables notes/.bin│ +│ (persisted to disk, one shard per note) │ └─────────────────────────────────────────────────────────────┘ ``` @@ -68,16 +68,18 @@ finding the right query becomes a significant productivity bottleneck. - **all-MiniLM-L6-v2**: 384-dimensional sentence embeddings - 86MB ONNX model (quantized version available at 22MB) -- Downloaded on first use to `zeppelin.search.index.path/models/` +- Installed before startup with `bin/install-search-model.sh` under + `zeppelin.search.index.path/models/` - Runs on CPU via ONNX Runtime (~5ms per paragraph) ### Index -- In-memory `ConcurrentHashMap` with `ReadWriteLock` +- In-memory `ConcurrentHashMap` - Each entry stores: embedding (384 floats), notebook name, paragraph text, title, extracted SQL table names, and paragraph output - 10K paragraphs ≈ 15MB RAM, 50K paragraphs ≈ 75MB RAM -- Persisted as versioned binary file (`embedding_index.bin`, currently v3) +- Persisted as one versioned binary shard per note under + `{zeppelin.search.index.path}/notes/`; updating a paragraph rewrites only its note's shard - Brute-force cosine similarity: < 50ms for 50K paragraphs ### What gets indexed (vs. LuceneSearch) @@ -126,8 +128,9 @@ Requires `zeppelin.search.enable = true` (already the default). ## Changes ### New files -- `zeppelin-server/.../search/EmbeddingSearch.java` — Core implementation (~700 lines) -- `zeppelin-server/.../search/EmbeddingSearchTest.java` — 11 tests including semantic validation +- `zeppelin-server/.../search/EmbeddingSearch.java` — Core implementation +- `zeppelin-server/.../search/EmbeddingSearchTest.java` — Tests including semantic + validation (requires the ONNX model, see Testing below) - `docs/embedding-search.md` — This document ### Modified files — Backend @@ -188,10 +191,12 @@ For Zeppelin's scale (typically < 50K paragraphs), brute-force cosine similarity on normalized vectors is fast enough (< 50ms), exact (no approximation error), and adds zero complexity. -### Why download model on first use instead of bundling? +### Why install the model ahead of time instead of bundling it? The ONNX model is 86MB. Bundling it would bloat the Zeppelin distribution. -Downloading on first use keeps the distribution lean and allows users to swap models. +Installing it explicitly with `bin/install-search-model.sh` keeps the distribution lean, +avoids network I/O and download delays during server startup, and works in air-gapped +environments once the pinned model artifacts have been staged. ### Why not use Lucene's vector search (since 9.0)? @@ -200,7 +205,8 @@ Zeppelin uses Lucene 8.7.0. Upgrading to 9.x is a separate, larger effort. ## Testing ```bash -# Run embedding search tests (requires model download, ~86MB first time) +# Install the pinned model once, then run embedding search tests +bin/install-search-model.sh ZEPPELIN_EMBEDDING_TEST=true mvn test -pl zeppelin-server \ -Dtest=EmbeddingSearchTest diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/search/EmbeddingSearch.java b/zeppelin-server/src/main/java/org/apache/zeppelin/search/EmbeddingSearch.java index 41b9666bd25..c3c33fb75ef 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/search/EmbeddingSearch.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/search/EmbeddingSearch.java @@ -27,15 +27,21 @@ import java.io.ByteArrayOutputStream; import java.io.DataInputStream; import java.io.DataOutputStream; +import java.io.File; import java.io.IOException; import java.nio.LongBuffer; +import java.nio.file.FileVisitResult; import java.nio.file.Files; import java.nio.file.Path; import java.nio.file.Paths; +import java.nio.file.SimpleFileVisitor; +import java.nio.file.StandardCopyOption; +import java.nio.file.attribute.BasicFileAttributes; import java.nio.file.attribute.PosixFilePermissions; import java.security.MessageDigest; import java.security.NoSuchAlgorithmException; import java.util.ArrayList; +import java.util.Iterator; import java.util.Locale; import java.util.Collections; import java.util.HashMap; @@ -48,11 +54,10 @@ import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicBoolean; -import java.util.concurrent.locks.ReadWriteLock; import java.util.concurrent.locks.ReentrantReadWriteLock; import java.util.regex.Matcher; import java.util.regex.Pattern; +import java.util.stream.Collectors; import javax.annotation.PreDestroy; import jakarta.inject.Inject; @@ -61,6 +66,7 @@ import org.apache.zeppelin.interpreter.InterpreterResult; import org.apache.zeppelin.interpreter.InterpreterResultMessage; import org.apache.zeppelin.notebook.Note; +import org.apache.zeppelin.notebook.NoteInfo; import org.apache.zeppelin.notebook.Notebook; import org.apache.zeppelin.notebook.Paragraph; import org.slf4j.Logger; @@ -74,12 +80,13 @@ * matched via cosine similarity, enabling natural language search like * "yesterday's spend query" to find {@code WHERE date = current_date - 1}. * - *

The embedding index is held in memory (float[][] + metadata) and persisted to a - * single binary file on disk. For typical Zeppelin deployments (< 50K paragraphs), - * brute-force cosine similarity completes in under 50ms. + *

The embedding index is held in memory (float[][] + metadata) and persisted to disk + * as one binary shard per note (see {@link #saveNoteShard(String)}), so editing a + * paragraph in one note only rewrites that note's shard. For typical Zeppelin + * deployments (< 50K paragraphs), brute-force cosine similarity completes in under 50ms. * - *

Model files are downloaded on first use to {@code zeppelin.search.index.path} - * and cached for subsequent starts. + *

Model files must be installed under {@code zeppelin.search.index.path} with + * {@code bin/install-search-model.sh} before semantic search is enabled. */ public class EmbeddingSearch extends SearchService { private static final Logger LOGGER = LoggerFactory.getLogger(EmbeddingSearch.class); @@ -136,9 +143,28 @@ public class EmbeddingSearch extends SearchService { * plausible deployment (~18 GB of vectors alone at 384 floats/entry). */ private static final int MAX_INDEX_ENTRIES = 10_000_000; - private static final String INDEX_FILE_NAME = "embedding_index.bin"; - /** Binary format version written by {@link #saveIndex()} and required by {@link #loadIndex()}. */ + /** Legacy single-file index, from before persistence was sharded by note (ZEPPELIN-6412). */ + private static final String LEGACY_INDEX_FILE_NAME = "embedding_index.bin"; + /** Binary format version of {@link #LEGACY_INDEX_FILE_NAME}, read only during migration. */ private static final int INDEX_VERSION = 3; + /** Subdirectory holding one binary shard per note. */ + private static final String NOTES_SHARD_DIR_NAME = "notes"; + private static final String SHARD_FILE_SUFFIX = ".bin"; + private static final String SHARD_TMP_SUFFIX = ".bin.tmp"; + /** + * Staging directory migration writes every note's shard into before anything touches + * {@link #NOTES_SHARD_DIR_NAME}. Published by a single atomic directory rename once every + * note has been staged successfully, so a reader can never observe a half-migrated + * {@code notes/} directory — see {@link #migrateLegacyIndexIfPresent()}. + */ + private static final String MIGRATION_STAGING_DIR_NAME = "notes.migrating"; + /** + * Binary format version written by {@link #saveNoteShard(String)} and required by + * {@link #loadNoteShard(String, Path)}. Independent of {@link #INDEX_VERSION}: these are + * different file formats at different paths, so conflating them would force migration + * code to special-case "right version, wrong location". + */ + private static final int SHARD_VERSION = 1; private static final String EXPECTED_MODEL_SHA256 = "6fd5d72fe4589f189f8ebc006442dbb529bb7ce38f8082112682524616046452"; @@ -152,8 +178,11 @@ public class EmbeddingSearch extends SearchService { // In-memory vector index: docId -> (embedding, metadata) private final ConcurrentHashMap index = new ConcurrentHashMap<>(); - private final ReadWriteLock indexLock = new ReentrantReadWriteLock(); - private final AtomicBoolean indexDirty = new AtomicBoolean(false); + /** One lock per note, created on demand; never removed (see {@link #lockFor(String)}). */ + private final ConcurrentHashMap noteLocks = + new ConcurrentHashMap<>(); + /** noteIds with in-memory changes not yet flushed to their shard file. */ + private final Set dirtyNoteIds = ConcurrentHashMap.newKeySet(); private final ScheduledExecutorService flushScheduler = Executors.newSingleThreadScheduledExecutor(r -> { Thread t = new Thread(r, "EmbeddingSearch-flush"); @@ -198,10 +227,13 @@ public EmbeddingSearch(ZeppelinConfiguration zConf, Notebook notebook) throws IO throw new IOException("Failed to initialize embedding model", e); } - boolean indexLoaded = loadIndex(); - if (shouldBootstrapIndex(zConf, indexLoaded)) { - notebook.addInitConsumer(this::addNoteIndex); - } + migrateLegacyIndexIfPresent(); + // Checked after migration (so a corrupt legacy file that migration deleted counts as + // "nothing persisted") but before loadAllShards() (so one individually-corrupt shard + // doesn't retroactively make this look like an empty-index deployment). + boolean hadPersistedIndex = shardsDirHasAnyShard(); + Set noteIdsToRebuild = loadAllShards(); + registerInitialRebuild(zConf, hadPersistedIndex, noteIdsToRebuild); flushScheduler.scheduleWithFixedDelay(this::flushIfDirty, FLUSH_INTERVAL_SECONDS, FLUSH_INTERVAL_SECONDS, TimeUnit.SECONDS); this.notebook.addNotebookEventListener(this); @@ -222,10 +254,10 @@ public EmbeddingSearch(ZeppelinConfiguration zConf, Notebook notebook) throws IO throw new IOException("Failed to initialize embedding model", e); } } - boolean indexLoaded = loadIndex(); - if (shouldBootstrapIndex(zConf, indexLoaded)) { - notebook.addInitConsumer(this::addNoteIndex); - } + migrateLegacyIndexIfPresent(); + boolean hadPersistedIndex = shardsDirHasAnyShard(); + Set noteIdsToRebuild = loadAllShards(); + registerInitialRebuild(zConf, hadPersistedIndex, noteIdsToRebuild); flushScheduler.scheduleWithFixedDelay(this::flushIfDirty, FLUSH_INTERVAL_SECONDS, FLUSH_INTERVAL_SECONDS, TimeUnit.SECONDS); this.notebook.addNotebookEventListener(this); @@ -247,6 +279,19 @@ private static void restrictPermissions(Path dir) { } } + /** Restrict a regular file (as opposed to a directory, see {@link #restrictPermissions}) to + * owner-only read/write, mirroring the permissions already applied to {@link #indexPath}. */ + private static void restrictFilePermissions(Path file) { + try { + if (Files.getFileStore(file).supportsFileAttributeView("posix")) { + Files.setPosixFilePermissions(file, + PosixFilePermissions.fromString("rw-------")); + } + } catch (IOException e) { + LOGGER.warn("Could not restrict permissions on {}", file, e); + } + } + // ---- Model initialization ---- private void initModel() throws OrtException, IOException { @@ -529,25 +574,25 @@ public List> query(String queryStr, Predicate readab String queryLower = queryStr.toLowerCase(Locale.ROOT); // Phase 1: find top-N results and discover relevant tables + // No lock is taken here: IndexEntry is immutable (all fields final) and + // ConcurrentHashMap gives weakly-consistent iteration — put()/remove() swap + // references atomically per key, so a concurrent write is seen either fully + // or not at all, never torn. Per-note locks (see lockFor) only need to protect + // the save path's per-note snapshot, not this read. List> scored = new ArrayList<>(); - indexLock.readLock().lock(); - try { - for (Map.Entry entry : index.entrySet()) { - // Dropping the entries here keeps them out of the table weights below and out of - // the cutoff, so the caller is served its own top results and not what is left of - // everyone's top results. - if (!readableNotes.computeIfAbsent(noteIdOf(entry.getKey()), readable::test)) { - continue; - } - float sim = cosineSimilarity(queryEmbedding, entry.getValue().embedding); - IndexEntry ie = entry.getValue(); - if (ie.text != null && ie.text.toLowerCase(Locale.ROOT).contains(queryLower)) { - sim += KEYWORD_BOOST; - } - scored.add(Map.entry(entry.getKey(), sim)); + for (Map.Entry entry : index.entrySet()) { + // Dropping the entries here keeps them out of the table weights below and out of + // the cutoff, so the caller is served its own top results and not what is left of + // everyone's top results. + if (!readableNotes.computeIfAbsent(noteIdOf(entry.getKey()), readable::test)) { + continue; } - } finally { - indexLock.readLock().unlock(); + float sim = cosineSimilarity(queryEmbedding, entry.getValue().embedding); + IndexEntry ie = entry.getValue(); + if (ie.text != null && ie.text.toLowerCase(Locale.ROOT).contains(queryLower)) { + sim += KEYWORD_BOOST; + } + scored.add(Map.entry(entry.getKey(), sim)); } scored.sort((a, b) -> Float.compare(b.getValue(), a.getValue())); @@ -633,7 +678,7 @@ public void addNoteIndex(String noteId) { } return null; }); - markDirty(); + markDirty(noteId); } catch (IOException e) { LOGGER.error("Failed to add note {} to index", noteId, e); } @@ -651,7 +696,7 @@ public void addParagraphIndex(String noteId, String paragraphId) { } return null; }); - markDirty(); + markDirty(noteId); } catch (IOException e) { LOGGER.error("Failed to add paragraph {} of note {}", paragraphId, noteId, e); } @@ -677,7 +722,7 @@ public void updateNoteIndex(String noteId) { if (newName == null) { return null; } - indexLock.writeLock().lock(); + lockFor(noteId).writeLock().lock(); try { boolean mutated = false; String notePrefix = noteId + "/"; @@ -695,10 +740,10 @@ public void updateNoteIndex(String noteId) { mutated = true; } if (mutated) { - markDirty(); + markDirty(noteId); } } finally { - indexLock.writeLock().unlock(); + lockFor(noteId).writeLock().unlock(); } return null; }); @@ -719,7 +764,7 @@ public void updateParagraphIndex(String noteId, String paragraphId) { } return null; }); - markDirty(); + markDirty(noteId); } catch (IOException e) { LOGGER.error("Failed to update paragraph {} of note {}", paragraphId, noteId, e); } @@ -730,14 +775,21 @@ public void deleteNoteIndex(String noteId) { if (noteId == null) { return; } - indexLock.writeLock().lock(); + lockFor(noteId).writeLock().lock(); try { index.entrySet().removeIf(e -> e.getKey().equals(noteId) || e.getKey().startsWith(noteId + "/")); } finally { - indexLock.writeLock().unlock(); + lockFor(noteId).writeLock().unlock(); } - markDirty(); + // Don't delete the shard file here. saveNoteShard() already deletes a note's shard + // when it finds no remaining entries for it, and it does that under this same note's + // lock for its *entire* save (snapshot + write) — so routing the delete through the + // normal dirty/flush path (instead of racing an out-of-band file delete against an + // in-flight save) is what keeps a concurrent flush from resurrecting a deleted note's + // shard with a stale pre-delete snapshot. A failed delete-on-flush is retried the same + // way any other failed flush is: flushIfDirty() re-marks the note dirty on IOException. + markDirty(noteId); } @Override @@ -748,8 +800,13 @@ public void deleteParagraphIndex(String noteId, String paragraphId) { String docId = paragraphId != null ? String.join("/", noteId, PARAGRAPH, paragraphId) : noteId; - index.remove(docId); - markDirty(); + lockFor(noteId).writeLock().lock(); + try { + index.remove(docId); + } finally { + lockFor(noteId).writeLock().unlock(); + } + markDirty(noteId); } @Override @@ -770,45 +827,84 @@ public void close() { } } - private void markDirty() { - indexDirty.set(true); + /** Lock guarding {@code index} entries belonging to {@code noteId}. Created on first use + * and never removed — noteIds aren't reused, and removing an entry while another thread + * might hold a reference to the same lock instance (from a concurrent computeIfAbsent) + * would risk two threads ending up on two different lock objects for the same note. */ + private ReentrantReadWriteLock lockFor(String noteId) { + return noteLocks.computeIfAbsent(noteId, k -> new ReentrantReadWriteLock()); + } + + private void markDirty(String noteId) { + dirtyNoteIds.add(noteId); } /** - * Decide whether to register the initial-indexing consumer. + * Decide whether to register the initial-indexing consumer for an operator-requested + * full rebuild. Per-note rebuilds triggered by a corrupt/missing shard are handled + * directly in {@link #loadAllShards()}, and a from-scratch bootstrap (no persisted index + * at all) is handled by the {@code hadPersistedIndex} check at the constructors' call + * sites — neither goes through this method. * - * @param zConf Zeppelin configuration (for {@code isIndexRebuild}) - * @param loaded whether {@link #loadIndex()} completed successfully - * @return {@code true} if the index needs to be (re)built from notebooks. Triggers when - * config requests rebuild, the index file is missing, or it was present but - * failed to load (corrupt/partial). A failed load also deletes the bad file so - * the rebuilt index is written fresh. + * @param zConf Zeppelin configuration (for {@code isIndexRebuild}) + * @return {@code true} if {@code zeppelin.search.index.rebuild} requests a full rebuild */ - private boolean shouldBootstrapIndex(ZeppelinConfiguration zConf, boolean loaded) { - Path indexFile = indexPath.resolve(INDEX_FILE_NAME); - boolean fileMissing = !Files.exists(indexFile); - boolean corrupt = !loaded; - if (corrupt && !fileMissing) { - try { - Files.deleteIfExists(indexFile); - LOGGER.warn("Deleted corrupt embedding index file {}; will rebuild", indexFile); - } catch (IOException e) { - LOGGER.warn("Failed to delete corrupt embedding index file {}; will rebuild anyway", - indexFile, e); - } + private boolean shouldBootstrapIndex(ZeppelinConfiguration zConf) { + return zConf.isIndexRebuild(); + } + + /** Register either a full bootstrap or the smallest per-note repair required at startup. */ + private void registerInitialRebuild(ZeppelinConfiguration zConf, boolean hadPersistedIndex, + Set noteIdsToRebuild) { + if (shouldBootstrapIndex(zConf) || !hadPersistedIndex) { + notebook.addInitConsumer(this::addNoteIndex); + } else if (!noteIdsToRebuild.isEmpty()) { + notebook.addInitConsumer(noteId -> { + if (noteIdsToRebuild.contains(noteId)) { + addNoteIndex(noteId); + } + }); + } + } + + /** + * @return {@code true} if the notes shard directory exists and holds at least one shard + * file. Used right after migration (and before {@link #loadAllShards()} can delete + * an individually-corrupt one) to tell "nothing has ever been persisted" — which + * should bootstrap the whole notebook, same as the old single-file code did when + * its one file was simply missing — apart from "something was persisted and a + * piece of it happens to be corrupt," which only needs a per-note rebuild. + */ + private boolean shardsDirHasAnyShard() { + Path notesDir = indexPath.resolve(NOTES_SHARD_DIR_NAME); + if (!Files.exists(notesDir)) { + return false; } - return zConf.isIndexRebuild() || fileMissing || corrupt; + File[] shardFiles = notesDir.toFile().listFiles((d, name) -> name.endsWith(SHARD_FILE_SUFFIX)); + return shardFiles != null && shardFiles.length > 0; } - private void flushIfDirty() { - if (indexDirty.compareAndSet(true, false)) { + /** + * Flush dirty note shards serially. The scheduled task and {@link #close()} can invoke this + * method concurrently; without synchronization both callers could drain the same weakly + * consistent dirty-set iterator and write/move the same {@code .bin.tmp} file at once. + */ + private synchronized void flushIfDirty() { + Set toFlush = new HashSet<>(); + Iterator it = dirtyNoteIds.iterator(); + while (it.hasNext()) { + toFlush.add(it.next()); + it.remove(); + } + for (String noteId : toFlush) { try { - saveIndex(); + saveNoteShard(noteId); } catch (IOException e) { - // Re-set dirty so the next scheduled tick retries the flush + // Re-mark dirty so the next scheduled tick retries the flush // instead of silently dropping the failed write until the next mutation. - indexDirty.set(true); - LOGGER.error("Failed to flush embedding index to disk; will retry on next tick", e); + dirtyNoteIds.add(noteId); + LOGGER.error("Failed to flush embedding shard for note {}; will retry on next tick", + noteId, e); } } } @@ -835,11 +931,11 @@ private void indexParagraph(String noteId, String noteName, Paragraph p) { String tables = extractTables(pText); String output = extractOutput(p); - indexLock.writeLock().lock(); + lockFor(noteId).writeLock().lock(); try { index.put(docId, new IndexEntry(emb, noteName, pText, title, tables, output)); } finally { - indexLock.writeLock().unlock(); + lockFor(noteId).writeLock().unlock(); } } @@ -852,91 +948,386 @@ static String formatId(String noteId, Paragraph p) { // ---- Persistence ---- + private Path shardFile(String noteId) { + return indexPath.resolve(NOTES_SHARD_DIR_NAME).resolve(noteId + SHARD_FILE_SUFFIX); + } + /** - * Save index to a binary file. - * Format: [int:version=INDEX_VERSION][int:count] then for each entry: - * [utf:docId] [utf:noteName] [utf:text] [utf:title] [utf:tables] [utf:output] [float[384]:embedding] + * Save one note's entries to its own binary shard. + * Format: [int:shardVersion=SHARD_VERSION][utf:noteId][int:count] then for each entry: + * [utf:docId] [utf:noteName] [utf:text] [utf:title] [utf:tables] [utf:output] + * [float[384]:embedding] + * + *

A single paragraph edit only rewrites the shard of the note it belongs to — + * every other note's shard is untouched (ZEPPELIN-6412). + * + *

The note's read lock is held for the entire save — snapshot and disk write + * alike — not just the snapshot. Releasing it in between would let a concurrent + * {@code deleteNoteIndex}/mutator run after the snapshot was taken but before the file + * write landed, so the write could recreate a shard the delete had just removed from + * memory (and expected removed on disk) with stale, pre-delete content. Holding the lock + * the whole time forces every mutator for this note to wait until the save is fully done, + * so a save always reflects a state that was current at some point and is never undone + * by a mutation it couldn't have known about. */ - // TODO(ZEPPELIN-6412): Shard persistence by note (e.g. index/notes/.bin) so a single - // paragraph edit only rewrites that note's file instead of the full index. Needs a per-note - // lock strategy, a manifest for load, and a compaction path for deletes; may also revisit - // append-only log + periodic compaction as the persistence model. - private void saveIndex() throws IOException { - Path file = indexPath.resolve(INDEX_FILE_NAME); - Path tmpFile = indexPath.resolve(INDEX_FILE_NAME + ".tmp"); - - // Serialize to buffer under lock - byte[] data; - indexLock.readLock().lock(); + private void saveNoteShard(String noteId) throws IOException { + lockFor(noteId).readLock().lock(); try { - ByteArrayOutputStream baos = new ByteArrayOutputStream(); - try (DataOutputStream out = new DataOutputStream(baos)) { - out.writeInt(INDEX_VERSION); - out.writeInt(index.size()); - for (Map.Entry e : index.entrySet()) { - out.writeUTF(e.getKey()); - out.writeUTF(e.getValue().noteName != null ? e.getValue().noteName : ""); - String text = e.getValue().text != null ? e.getValue().text : ""; - if (text.length() > MAX_PERSISTED_TEXT_LENGTH) { - text = text.substring(0, MAX_PERSISTED_TEXT_LENGTH); - } - out.writeUTF(text); - out.writeUTF(e.getValue().title != null ? e.getValue().title : ""); - out.writeUTF(e.getValue().tables != null ? e.getValue().tables : ""); - String output = e.getValue().output != null ? e.getValue().output : ""; - if (output.length() > MAX_PERSISTED_OUTPUT_LENGTH) { - output = output.substring(0, MAX_PERSISTED_OUTPUT_LENGTH); - } - out.writeUTF(output); - for (float v : e.getValue().embedding) { - out.writeFloat(v); - } + List> entries = index.entrySet().stream() + .filter(e -> e.getKey().equals(noteId) || e.getKey().startsWith(noteId + "/")) + .collect(Collectors.toList()); + + Path file = shardFile(noteId); + if (entries.isEmpty()) { + // No entries left for this note (e.g. deleteNoteIndex ran first) — remove the shard + // rather than leave an empty file. If this delete itself fails, the IOException + // propagates out and flushIfDirty() re-marks the note dirty to retry. + Files.deleteIfExists(file); + return; + } + + byte[] data = serializeShard(noteId, entries); + + Files.createDirectories(indexPath.resolve(NOTES_SHARD_DIR_NAME)); + Path tmpFile = indexPath.resolve(NOTES_SHARD_DIR_NAME).resolve(noteId + SHARD_TMP_SUFFIX); + Files.write(tmpFile, data); + Files.move(tmpFile, file, StandardCopyOption.REPLACE_EXISTING, + StandardCopyOption.ATOMIC_MOVE); + restrictFilePermissions(file); + } finally { + lockFor(noteId).readLock().unlock(); + } + } + + private static byte[] serializeShard(String noteId, List> entries) + throws IOException { + ByteArrayOutputStream baos = new ByteArrayOutputStream(); + try (DataOutputStream out = new DataOutputStream(baos)) { + out.writeInt(SHARD_VERSION); + out.writeUTF(noteId); + out.writeInt(entries.size()); + for (Map.Entry e : entries) { + out.writeUTF(e.getKey()); + out.writeUTF(e.getValue().noteName != null ? e.getValue().noteName : ""); + String text = e.getValue().text != null ? e.getValue().text : ""; + if (text.length() > MAX_PERSISTED_TEXT_LENGTH) { + text = text.substring(0, MAX_PERSISTED_TEXT_LENGTH); + } + out.writeUTF(text); + out.writeUTF(e.getValue().title != null ? e.getValue().title : ""); + out.writeUTF(e.getValue().tables != null ? e.getValue().tables : ""); + String output = e.getValue().output != null ? e.getValue().output : ""; + if (output.length() > MAX_PERSISTED_OUTPUT_LENGTH) { + output = output.substring(0, MAX_PERSISTED_OUTPUT_LENGTH); + } + out.writeUTF(output); + for (float v : e.getValue().embedding) { + out.writeFloat(v); + } + } + } + return baos.toByteArray(); + } + + /** + * Load every note's shard from {@code {indexPath}/notes/}. There's no manifest file — + * like {@code LocalRecoveryStorage}, the shard directory is scanned directly and each + * note's id is recovered from its filename. + * + * @return IDs of notes whose shard is corrupt or missing. A corrupt shard does not fail the + * whole load: it's deleted and returned for per-note rebuilding. Likewise, comparing + * the persisted shard IDs with {@link Notebook#getNotesInfo()} lets a single vanished + * shard self-heal without forcing a full rebuild. + */ + private Set loadAllShards() { + Set noteIdsToRebuild = new HashSet<>(); + Set currentNoteIds = notebook.getNotesInfo().stream() + .map(NoteInfo::getId) + .collect(Collectors.toSet()); + Path notesDir = indexPath.resolve(NOTES_SHARD_DIR_NAME); + if (!Files.exists(notesDir)) { + return noteIdsToRebuild; + } + File[] shardFiles = notesDir.toFile().listFiles((d, name) -> name.endsWith(SHARD_FILE_SUFFIX)); + if (shardFiles == null) { + return noteIdsToRebuild; + } + Set persistedNoteIds = new HashSet<>(); + for (File f : shardFiles) { + String noteId = f.getName().substring(0, f.getName().length() - SHARD_FILE_SUFFIX.length()); + if (!currentNoteIds.contains(noteId)) { + LOGGER.info("Deleting orphan embedding shard {} for note {} that no longer exists", + f, noteId); + try { + Files.deleteIfExists(f.toPath()); + } catch (IOException e) { + // Do not load stale entries into memory. Route the failed file deletion through + // the normal dirty flush path so the scheduler retries it after startup. + markDirty(noteId); + LOGGER.warn("Failed to delete orphan shard {}; will retry on next flush", f, e); + } + continue; + } + persistedNoteIds.add(noteId); + if (!loadNoteShard(noteId, f.toPath())) { + LOGGER.warn("Corrupt embedding shard for note {} ({}); deleting and rebuilding " + + "just that note", noteId, f); + try { + Files.deleteIfExists(f.toPath()); + } catch (IOException e) { + LOGGER.warn("Failed to delete corrupt shard {}; will attempt rebuild anyway", f, e); + } + noteIdsToRebuild.add(noteId); + } + } + + for (String noteId : currentNoteIds) { + if (!persistedNoteIds.contains(noteId)) { + LOGGER.warn("Embedding shard for note {} is missing; rebuilding just that note", + noteId); + noteIdsToRebuild.add(noteId); + } + } + return noteIdsToRebuild; + } + + /** + * Load a single note's shard into {@link #index}. Returns {@code false} on any + * corruption (bad version, filename/content noteId mismatch, bad count, an entry whose + * docId doesn't belong to this note, or truncated/malformed data partway through), + * signalling the caller to discard the file and rebuild just this note. + * + *

Every entry is read into a local map first; {@link #index} is only touched once, at + * the end, via {@link #commitNoteEntries}, and only if every entry the header promised + * was read and validated successfully. A shard that fails partway through — even after + * successfully reading one or more entries — contributes nothing: there's no path where + * some of a failed shard's entries end up live in {@link #index} while others don't. + */ + private boolean loadNoteShard(String noteId, Path file) { + Map loaded = new HashMap<>(); + try (DataInputStream in = new DataInputStream(Files.newInputStream(file))) { + int version = in.readInt(); + if (version != SHARD_VERSION) { + LOGGER.warn("Shard {} version {} does not match expected {}", file, version, + SHARD_VERSION); + return false; + } + String headerNoteId = in.readUTF(); + if (!headerNoteId.equals(noteId)) { + LOGGER.warn("Shard {} noteId {} does not match filename-derived noteId {}", + file, headerNoteId, noteId); + return false; + } + int count = in.readInt(); + if (count < 0 || count > MAX_INDEX_ENTRIES) { + LOGGER.error("Shard {} entry count {} exceeds sanity bound ({})", file, count, + MAX_INDEX_ENTRIES); + return false; + } + for (int i = 0; i < count; i++) { + String docId = in.readUTF(); + if (!docId.equals(noteId) && !docId.startsWith(noteId + "/")) { + LOGGER.warn("Shard {} entry {} has a docId that doesn't belong to note {}; " + + "treating the whole shard as corrupt", file, docId, noteId); + return false; + } + String noteName = in.readUTF(); + String text = in.readUTF(); + String title = in.readUTF(); + String tables = in.readUTF(); + String output = in.readUTF(); + float[] emb = new float[EMBEDDING_DIM]; + for (int j = 0; j < EMBEDDING_DIM; j++) { + emb[j] = in.readFloat(); } + loaded.put(docId, new IndexEntry(emb, noteName, text, title, tables, output)); } - data = baos.toByteArray(); + } catch (IOException e) { + LOGGER.warn("Failed to load embedding shard {}; discarding {} entry/entries read " + + "before the failure", file, loaded.size(), e); + return false; + } + commitNoteEntries(noteId, loaded); + return true; + } + + /** Atomically replace this note's entries in {@link #index} with exactly {@code entries} — + * the commit point a successful {@link #loadNoteShard} uses, so a shard's entries become + * visible all at once or not at all. */ + private void commitNoteEntries(String noteId, Map entries) { + lockFor(noteId).writeLock().lock(); + try { + index.entrySet().removeIf(e -> + e.getKey().equals(noteId) || e.getKey().startsWith(noteId + "/")); + index.putAll(entries); } finally { - indexLock.readLock().unlock(); + lockFor(noteId).writeLock().unlock(); + } + } + + /** + * One-time migration from the pre-ZEPPELIN-6412 single-file index. If + * {@code embedding_index.bin} is present, load it with the legacy format, stage every + * note's shard in {@link #MIGRATION_STAGING_DIR_NAME}, and only if every single one of + * them is staged successfully, publish the whole staging directory as {@link + * #NOTES_SHARD_DIR_NAME} with one atomic rename — then delete the legacy file so this + * never runs again. + * + *

This is all-or-nothing on purpose: writing shards straight into the real {@code + * notes/} directory one at a time would let a mid-migration failure (disk full, a bad + * noteId, anything) leave some notes migrated and others not, and {@link + * #loadAllShards()} has no way to tell that apart from a complete, correct index — it + * would just load the partial set and serve it. Staging first and publishing with a + * single rename means a reader only ever sees "no {@code notes/} directory yet" or "a + * fully-migrated one," never something in between. If staging fails, the legacy file is + * left in place (so the next restart retries from scratch; shard writes are keyed by + * noteId and idempotent, so repeating this is safe) and the staging directory is cleared. + */ + private void migrateLegacyIndexIfPresent() { + Path legacyFile = indexPath.resolve(LEGACY_INDEX_FILE_NAME); + if (!Files.exists(legacyFile)) { + return; + } + // A published shard set is the newer format and therefore the source of truth. This state + // can legitimately occur when a previous migration published notes/ successfully but failed + // to delete the legacy file. Re-running migration here would replace potentially newer shards + // with the stale legacy snapshot on every restart. + if (shardsDirHasAnyShard()) { + LOGGER.info("Per-note embedding shards already exist; ignoring and deleting stale legacy " + + "index {}", legacyFile); + deleteLegacyFile(legacyFile); + return; + } + LOGGER.info("Migrating legacy single-file embedding index {} to per-note shards", + legacyFile); + if (!loadLegacyIndex(legacyFile)) { + LOGGER.warn("Legacy index {} failed to load during migration; deleting and " + + "bootstrapping a fresh index", legacyFile); + deleteLegacyFile(legacyFile); + index.clear(); + return; } - // Write to disk outside lock - Files.write(tmpFile, data); - Files.move(tmpFile, file, java.nio.file.StandardCopyOption.REPLACE_EXISTING, - java.nio.file.StandardCopyOption.ATOMIC_MOVE); - // Restrict file permissions + Path stagingDir = indexPath.resolve(MIGRATION_STAGING_DIR_NAME); + Path notesDir = indexPath.resolve(NOTES_SHARD_DIR_NAME); try { - if (Files.getFileStore(file).supportsFileAttributeView("posix")) { - Files.setPosixFilePermissions(file, - PosixFilePermissions.fromString("rw-------")); + // Clear any half-written staging debris left by a migration attempt that crashed + // (as opposed to failing cleanly through the catch block below) before this one. + deleteDirectoryRecursively(stagingDir); + Files.createDirectories(stagingDir); + + Set noteIds = index.keySet().stream() + .map(SearchService::noteIdOf) + .collect(Collectors.toSet()); + for (String noteId : noteIds) { + writeShardToDir(stagingDir, noteId); + } + + // Publish: this is the only point that touches the real notes/ directory, and it's a + // single rename — so a reader can never observe a partially-migrated one. + if (Files.exists(notesDir)) { + // Not expected (migration always runs before loadAllShards() creates/uses notes/), + // but don't let a stray pre-existing directory fail the whole migration. + deleteDirectoryRecursively(notesDir); } + Files.move(stagingDir, notesDir, StandardCopyOption.ATOMIC_MOVE); + + LOGGER.info("Migrated {} notes ({} entries) to per-note shards", noteIds.size(), + index.size()); } catch (IOException e) { - LOGGER.warn("Could not restrict permissions on {}", file, e); + LOGGER.error("Failed to stage per-note shards during migration; legacy index kept " + + "for retry on next restart", e); + try { + deleteDirectoryRecursively(stagingDir); + } catch (IOException cleanup) { + LOGGER.warn("Failed to clean up migration staging directory {}", stagingDir, cleanup); + } + index.clear(); + return; + } + // index.clear() here (success path only) so loadAllShards() is the one source of truth + // for what ends up in memory, instead of trusting this migration pass's copy. + index.clear(); + deleteLegacyFile(legacyFile); + } + + /** Write one note's entries, if any, as a shard file directly under {@code dir} — used by + * migration to stage into {@link #MIGRATION_STAGING_DIR_NAME}, never the real shard + * directory. Unlike {@link #saveNoteShard}, there's no existing file to delete when a + * note has no entries: the staging directory starts empty, so there's simply nothing to + * write for it. No note-level lock is needed here — migration runs synchronously in the + * constructor, before {@code notebook.addNotebookEventListener}/{@code addInitConsumer} + * are registered, so nothing else can be mutating {@link #index} concurrently yet. */ + private void writeShardToDir(Path dir, String noteId) throws IOException { + List> entries = index.entrySet().stream() + .filter(e -> e.getKey().equals(noteId) || e.getKey().startsWith(noteId + "/")) + .collect(Collectors.toList()); + if (entries.isEmpty()) { + return; + } + byte[] data = serializeShard(noteId, entries); + Path file = dir.resolve(noteId + SHARD_FILE_SUFFIX); + Files.write(file, data); + restrictFilePermissions(file); + } + + /** Recursively delete {@code dir} if it exists. Used to clear migration staging debris — + * never called on {@link #NOTES_SHARD_DIR_NAME} except right before it's about to be + * replaced by a freshly-published staging directory. */ + private static void deleteDirectoryRecursively(Path dir) throws IOException { + if (!Files.exists(dir)) { + return; + } + Files.walkFileTree(dir, new SimpleFileVisitor() { + @Override + public FileVisitResult visitFile(Path file, BasicFileAttributes attrs) throws IOException { + // Propagate the first failure so migration cannot publish a directory containing + // debris from an earlier attempt. + Files.delete(file); + return FileVisitResult.CONTINUE; + } + + @Override + public FileVisitResult postVisitDirectory(Path directory, IOException error) + throws IOException { + if (error != null) { + throw error; + } + Files.delete(directory); + return FileVisitResult.CONTINUE; + } + }); + } + + private static void deleteLegacyFile(Path legacyFile) { + try { + Files.delete(legacyFile); + } catch (IOException e) { + LOGGER.warn("Failed to delete legacy embedding index {} after migration", legacyFile, e); } } /** - * Load the index from disk. + * Read the legacy (pre-ZEPPELIN-6412) single-file index format. Used only by + * {@link #migrateLegacyIndexIfPresent()} — once that file is deleted, this never runs + * again. * - * @return {@code true} if the index loaded successfully (or file was absent); - * {@code false} if the file was present but failed to load or was corrupt, - * signalling the caller to trigger a bootstrap rebuild. + * @return {@code true} if the file loaded successfully; {@code false} if corrupt. */ - private boolean loadIndex() { - Path file = indexPath.resolve(INDEX_FILE_NAME); - if (!Files.exists(file)) { - return true; - } + private boolean loadLegacyIndex(Path file) { try (DataInputStream in = new DataInputStream(Files.newInputStream(file))) { int version = in.readInt(); if (version != INDEX_VERSION) { - LOGGER.warn("Index file version {} does not match expected {}; treating as corrupt " - + "and rebuilding", version, INDEX_VERSION); + LOGGER.warn("Legacy index file version {} does not match expected {}; treating as " + + "corrupt", version, INDEX_VERSION); return false; } int count = in.readInt(); - LOGGER.info("Loading {} embedding index entries (v{}) from {}", count, version, file); + LOGGER.info("Loading {} embedding index entries (v{}) from legacy file {}", count, + version, file); if (count < 0 || count > MAX_INDEX_ENTRIES) { - LOGGER.error("Index entry count {} exceeds sanity bound ({}), treating as corrupt", - count, MAX_INDEX_ENTRIES); + LOGGER.error("Legacy index entry count {} exceeds sanity bound ({}), treating as " + + "corrupt", count, MAX_INDEX_ENTRIES); return false; } for (int i = 0; i < count; i++) { @@ -952,11 +1343,10 @@ private boolean loadIndex() { } index.put(docId, new IndexEntry(emb, noteName, text, title, tables, output)); } - LOGGER.info("Loaded {} entries into embedding index", index.size()); + LOGGER.info("Loaded {} entries from legacy index", index.size()); return true; } catch (IOException e) { - LOGGER.warn("Failed to load embedding index from {}; will rebuild on init", file, e); - // Clear any partially-loaded state so we start from a clean slate on rebuild. + LOGGER.warn("Failed to load legacy embedding index from {}", file, e); index.clear(); return false; } diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/search/EmbeddingSearchShardingTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/search/EmbeddingSearchShardingTest.java new file mode 100644 index 00000000000..d60ff38428c --- /dev/null +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/search/EmbeddingSearchShardingTest.java @@ -0,0 +1,786 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.zeppelin.search; + +import static org.junit.jupiter.api.Assertions.assertArrayEquals; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.api.Assumptions.assumeFalse; +import static org.junit.jupiter.api.Assumptions.assumeTrue; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import java.io.DataOutputStream; +import java.io.File; +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.attribute.PosixFilePermissions; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; + +import org.apache.commons.io.FileUtils; +import org.apache.zeppelin.conf.ZeppelinConfiguration; +import org.apache.zeppelin.interpreter.InterpreterFactory; +import org.apache.zeppelin.interpreter.InterpreterSetting; +import org.apache.zeppelin.interpreter.InterpreterSettingManager; +import org.apache.zeppelin.notebook.AuthorizationService; +import org.apache.zeppelin.notebook.NoteManager; +import org.apache.zeppelin.notebook.Notebook; +import org.apache.zeppelin.notebook.Paragraph; +import org.apache.zeppelin.notebook.repo.InMemoryNotebookRepo; +import org.apache.zeppelin.notebook.repo.NotebookRepo; +import org.apache.zeppelin.user.AuthenticationInfo; +import org.apache.zeppelin.user.Credentials; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +/** + * Tests for ZEPPELIN-6412: {@link EmbeddingSearch} persistence sharded by note. + * + *

Unlike {@link EmbeddingSearchTest}, these use the package-private {@code skipModel} + * constructor, so they don't need the ONNX model and aren't gated behind + * {@code ZEPPELIN_EMBEDDING_TEST} — {@code embed()} falls back to a zero vector when + * there's no model, which doesn't matter to the persistence format under test here. + */ +class EmbeddingSearchShardingTest { + + private static final String NOTES_DIR = "notes"; + private static final String SHARD_SUFFIX = ".bin"; + private static final String LEGACY_FILE = "embedding_index.bin"; + + private File indexDir; + private ZeppelinConfiguration zConf; + private Notebook notebook; + + @BeforeEach + void startUp() throws IOException { + indexDir = Files.createTempDirectory(this.getClass().getSimpleName()).toFile(); + zConf = ZeppelinConfiguration.load(); + zConf.setProperty(ZeppelinConfiguration.ConfVars.ZEPPELIN_SEARCH_INDEX_PATH.getVarName(), + indexDir.getAbsolutePath()); + + NoteManager noteManager = new NoteManager(new InMemoryNotebookRepo(), zConf); + InterpreterSettingManager interpreterSettingManager = mock(InterpreterSettingManager.class); + InterpreterSetting defaultInterpreterSetting = mock(InterpreterSetting.class); + when(defaultInterpreterSetting.getName()).thenReturn("test"); + when(interpreterSettingManager.getDefaultInterpreterSetting()) + .thenReturn(defaultInterpreterSetting); + notebook = new Notebook(zConf, mock(AuthorizationService.class), + mock(NotebookRepo.class), noteManager, + mock(InterpreterFactory.class), interpreterSettingManager, + mock(Credentials.class), null); + } + + @AfterEach + void shutDown() throws IOException { + FileUtils.deleteDirectory(indexDir); + } + + private EmbeddingSearch openSearch() throws IOException { + return new EmbeddingSearch(zConf, notebook, true); + } + + private File shardFile(String noteId) { + return new File(new File(indexDir, NOTES_DIR), noteId + SHARD_SUFFIX); + } + + private String newNoteWithParagraph(String noteName, String text) throws IOException { + String noteId = notebook.createNote(noteName, AuthenticationInfo.ANONYMOUS); + notebook.processNote(noteId, note -> { + Paragraph p = note.addNewParagraph(AuthenticationInfo.ANONYMOUS); + p.setText(text); + return null; + }); + return noteId; + } + + private String lastParagraphId(String noteId) throws IOException { + AtomicReference id = new AtomicReference<>(); + notebook.processNote(noteId, note -> { + id.set(note.getLastParagraph().getId()); + return null; + }); + return id.get(); + } + + @Test + void savesOneShardFilePerNote() throws IOException { + String noteA = newNoteWithParagraph("NoteA", "select * from a"); + String noteB = newNoteWithParagraph("NoteB", "select * from b"); + + EmbeddingSearch search = openSearch(); + search.addNoteIndex(noteA); + search.addNoteIndex(noteB); + search.close(); + + assertTrue(shardFile(noteA).exists(), "note A should have its own shard file"); + assertTrue(shardFile(noteB).exists(), "note B should have its own shard file"); + File[] shards = new File(indexDir, NOTES_DIR).listFiles((d, n) -> n.endsWith(SHARD_SUFFIX)); + assertEquals(2, shards.length, "exactly one shard per note, no shared file"); + } + + @Test + void editingOneNoteDoesNotTouchOtherNotesShardFile() throws IOException { + String noteA = newNoteWithParagraph("NoteA", "select * from a"); + String noteB = newNoteWithParagraph("NoteB", "select * from b"); + + EmbeddingSearch search = openSearch(); + search.addNoteIndex(noteA); + search.addNoteIndex(noteB); + search.close(); // first flush: both shards written + + byte[] noteBBytesBefore = Files.readAllBytes(shardFile(noteB).toPath()); + byte[] noteABytesBefore = Files.readAllBytes(shardFile(noteA).toPath()); + + // Reopen (loads both shards) and edit only note A. + EmbeddingSearch search2 = openSearch(); + String paragraphId = lastParagraphId(noteA); + notebook.processNote(noteA, note -> { + note.getLastParagraph().setText("select * from a_renamed_table"); + return null; + }); + search2.updateParagraphIndex(noteA, paragraphId); + search2.close(); // second flush: only note A changed + + byte[] noteABytesAfter = Files.readAllBytes(shardFile(noteA).toPath()); + byte[] noteBBytesAfter = Files.readAllBytes(shardFile(noteB).toPath()); + + assertArrayEquals(noteBBytesBefore, noteBBytesAfter, + "untouched note's shard must be byte-for-byte identical — this is the ZEPPELIN-6412 fix"); + assertFalse(java.util.Arrays.equals(noteABytesBefore, noteABytesAfter), + "edited note's shard should have actually changed"); + } + + @Test + void multipleParagraphEditsInSameNoteAreReflectedAfterOneFlush() throws IOException { + String noteId = notebook.createNote("MultiPara", AuthenticationInfo.ANONYMOUS); + List paragraphIds = new ArrayList<>(); + notebook.processNote(noteId, note -> { + for (int i = 0; i < 2; i++) { + Paragraph p = note.addNewParagraph(AuthenticationInfo.ANONYMOUS); + p.setText("original text " + i); + paragraphIds.add(p.getId()); + } + return null; + }); + + EmbeddingSearch search = openSearch(); + search.addNoteIndex(noteId); + search.close(); + + // Edit both paragraphs before the next flush — should collapse into one shard write. + EmbeddingSearch search2 = openSearch(); + notebook.processNote(noteId, note -> { + note.getParagraph(paragraphIds.get(0)).setText("updated alpha content"); + note.getParagraph(paragraphIds.get(1)).setText("updated beta content"); + return null; + }); + search2.updateParagraphIndex(noteId, paragraphIds.get(0)); + search2.updateParagraphIndex(noteId, paragraphIds.get(1)); + search2.close(); + + EmbeddingSearch search3 = openSearch(); + List> alphaResults = search3.query("alpha content", id -> true); + List> betaResults = search3.query("beta content", id -> true); + search3.close(); + + assertTrue(alphaResults.stream().anyMatch(r -> r.get("text").contains("alpha")), + "first paragraph's edit should be in the single flushed shard"); + assertTrue(betaResults.stream().anyMatch(r -> r.get("text").contains("beta")), + "second paragraph's edit should be in the same single flushed shard"); + } + + @Test + void deletingNoteRemovesItsShardButNotOthers() throws IOException { + String noteA = newNoteWithParagraph("NoteA", "select * from a"); + String noteB = newNoteWithParagraph("NoteB", "select * from b"); + + EmbeddingSearch search = openSearch(); + search.addNoteIndex(noteA); + search.addNoteIndex(noteB); + search.close(); + assertTrue(shardFile(noteA).exists()); + assertTrue(shardFile(noteB).exists()); + + EmbeddingSearch search2 = openSearch(); + search2.deleteNoteIndex(noteA); + // deleteNoteIndex() only mutates memory and marks the note dirty now — the shard file + // removal happens on the next flush, under this note's lock, so it can never race with + // an in-flight save that captured a snapshot before the delete (see saveNoteShard's + // javadoc). So the file is still there until that flush happens. + assertTrue(shardFile(noteA).exists(), + "shard removal is deferred to the next flush, not done inline"); + search2.close(); // the deferred flush actually removes it + + assertFalse(shardFile(noteA).exists(), "deleted note's shard should be gone after a flush"); + assertTrue(shardFile(noteB).exists(), "other note's shard must survive"); + } + + @Test + void deletingParagraphsShrinksThenRemovesTheShard() throws IOException { + String noteId = notebook.createNote("TwoParas", AuthenticationInfo.ANONYMOUS); + List paragraphIds = new ArrayList<>(); + notebook.processNote(noteId, note -> { + for (int i = 0; i < 2; i++) { + Paragraph p = note.addNewParagraph(AuthenticationInfo.ANONYMOUS); + p.setText("paragraph number " + i); + paragraphIds.add(p.getId()); + } + return null; + }); + + EmbeddingSearch search = openSearch(); + search.addNoteIndex(noteId); + search.close(); + assertTrue(shardFile(noteId).exists()); + + // Delete one of two paragraphs: shard should remain, with only the other entry. + EmbeddingSearch search2 = openSearch(); + search2.deleteParagraphIndex(noteId, paragraphIds.get(0)); + search2.close(); + assertTrue(shardFile(noteId).exists(), "shard should remain while one paragraph is left"); + + EmbeddingSearch search3 = openSearch(); + List> remaining = search3.query("paragraph number 1", id -> true); + assertTrue(remaining.stream().anyMatch(r -> r.get("text").contains("paragraph number 1"))); + + // Delete the last paragraph: the shard file should disappear entirely. + search3.deleteParagraphIndex(noteId, paragraphIds.get(1)); + search3.close(); + assertFalse(shardFile(noteId).exists(), + "shard with no remaining entries should be deleted, not left empty"); + } + + @Test + void migrationFromLegacyIndexFileProducesShardsAndDeletesLegacyFile() throws IOException { + String noteId = newNoteWithParagraph("Legacy Title", "current notebook content"); + String docId = noteId + "/paragraph/p1"; + writeLegacyIndexFile(new File(indexDir, LEGACY_FILE), + new LegacyEntry(docId, noteId, "legacy migrated content about quarterly revenue", + "Legacy Title")); + + EmbeddingSearch search = openSearch(); // constructor runs migration before loading shards + + assertFalse(new File(indexDir, LEGACY_FILE).exists(), + "legacy single-file index should be deleted after migration"); + assertTrue(shardFile(noteId).exists(), "migration should have written a shard for the note"); + + List> results = search.query("quarterly revenue", id -> true); + assertTrue(results.stream().anyMatch(r -> r.get("text").contains("quarterly revenue")), + "migrated content should be queryable without any further indexing"); + search.close(); + } + + @Test + void existingShardsTakePrecedenceOverStaleLegacyIndex() throws IOException { + String noteId = newNoteWithParagraph("CurrentNote", "current shard content"); + + EmbeddingSearch initial = openSearch(); + initial.addNoteIndex(noteId); + initial.close(); + byte[] shardBefore = Files.readAllBytes(shardFile(noteId).toPath()); + + File legacyFile = new File(indexDir, LEGACY_FILE); + writeLegacyIndexFile(legacyFile, + new LegacyEntry(noteId + "/paragraph/legacy", "CurrentNote", + "stale legacy content", "Stale")); + + EmbeddingSearch reopened = openSearch(); + + assertFalse(legacyFile.exists(), + "a stale legacy file should be discarded when published shards already exist"); + assertArrayEquals(shardBefore, Files.readAllBytes(shardFile(noteId).toPath()), + "reopening must not replace newer shards with the stale legacy snapshot"); + assertTrue(reopened.query("current shard content", id -> true).stream() + .anyMatch(r -> r.get("text").contains("current shard content"))); + assertTrue(reopened.query("stale legacy content", id -> true).isEmpty(), + "stale legacy entries must not become searchable again"); + reopened.close(); + } + + @Test + void concurrentCloseFlushesLeaveReadableShards() throws Exception { + String noteA = newNoteWithParagraph("ConcurrentFlushA", "concurrent flush alpha"); + String noteB = newNoteWithParagraph("ConcurrentFlushB", "concurrent flush beta"); + EmbeddingSearch search = openSearch(); + search.addNoteIndex(noteA); + search.addNoteIndex(noteB); + + CountDownLatch start = new CountDownLatch(1); + CountDownLatch done = new CountDownLatch(2); + AtomicReference failure = new AtomicReference<>(); + Runnable close = () -> { + try { + start.await(); + search.close(); + } catch (Throwable t) { + failure.compareAndSet(null, t); + } finally { + done.countDown(); + } + }; + new Thread(close).start(); + new Thread(close).start(); + start.countDown(); + assertTrue(done.await(10, TimeUnit.SECONDS)); + assertNull(failure.get(), "concurrent lifecycle flushes must not race on shard temp files"); + + EmbeddingSearch verify = openSearch(); + assertTrue(verify.query("concurrent flush alpha", id -> true).stream() + .anyMatch(r -> r.get("text").contains("concurrent flush alpha"))); + assertTrue(verify.query("concurrent flush beta", id -> true).stream() + .anyMatch(r -> r.get("text").contains("concurrent flush beta"))); + verify.close(); + } + + @Test + void bootstrapsExistingNotesWhenNoIndexExists() throws IOException, InterruptedException { + // No legacy file, no notes/ directory, no shards at all — just a Notebook with content + // that was never indexed. shouldBootstrapIndex() alone used to only check + // zeppelin.search.index.rebuild, so this notebook would silently stay unsearchable. + String noteId = newNoteWithParagraph("NeverIndexed", + "select * from completely_unindexed_table"); + + EmbeddingSearch search = openSearch(); // constructor should register a full bootstrap + + notebook.initNotebook(); + assertTrue(notebook.waitForFinishInit(10, TimeUnit.SECONDS)); + + List> results = search.query("completely_unindexed_table", id -> true); + assertTrue(results.stream().anyMatch(r -> r.get("text").contains("completely_unindexed_table")), + "an existing note should be indexed on startup when there's no persisted index at all"); + search.close(); // force the flush that actually persists what the bootstrap indexed + assertTrue(shardFile(noteId).exists(), "the bootstrap should have flushed a shard for it"); + } + + @Test + void corruptLegacyIndexFileAlsoTriggersFullBootstrap() throws IOException, InterruptedException { + // A legacy file that exists but fails to load (bad version) should be treated the same + // as "no index at all" once migration deletes it — not silently leave the notebook + // unindexed just because something used to be there. + String noteId = newNoteWithParagraph("WasNeverMigrated", + "select * from table_behind_corrupt_legacy_index"); + File legacyFile = new File(indexDir, LEGACY_FILE); + Files.write(legacyFile.toPath(), new byte[] {0, 0, 0, 99, 1, 2, 3}); // bad version header + + EmbeddingSearch search = openSearch(); + assertFalse(legacyFile.exists(), "unreadable legacy file should be deleted during migration"); + + notebook.initNotebook(); + assertTrue(notebook.waitForFinishInit(10, TimeUnit.SECONDS)); + + List> results = + search.query("table_behind_corrupt_legacy_index", id -> true); + assertTrue(results.stream() + .anyMatch(r -> r.get("text").contains("table_behind_corrupt_legacy_index")), + "a corrupt legacy file must not leave the notebook unindexed"); + search.close(); + } + + @Test + void truncatedShardDoesNotLeavePartiallyLoadedEntries() throws IOException { + String noteId = newNoteWithParagraph("TruncatedNote", "real notebook content"); + String survivingText = "truncated shard entry alpha content that must not leak"; + File notesDir = new File(indexDir, NOTES_DIR); + assertTrue(notesDir.mkdirs() || notesDir.isDirectory()); + writeTruncatedShard(shardFile(noteId), noteId, survivingText); + + EmbeddingSearch search = openSearch(); // loadNoteShard() should fail cleanly on this shard + + List> results = search.query(survivingText, id -> true); + assertTrue(results.isEmpty(), + "a shard that fails partway through must not leave its earlier, successfully-read " + + "entries live in the index"); + assertFalse(shardFile(noteId).exists(), "the corrupt shard file should be deleted"); + search.close(); + } + + @Test + void concurrentFlushCannotRecreateDeletedNoteShard() throws Exception { + // Create every note up front, before any EmbeddingSearch (and so any notebook event + // listener) exists. Each iteration below opens and closes its own instances, and a + // note-create event fired *after* an earlier iteration's closed instance is still + // registered as a listener would hit a RejectedExecutionException unrelated to the + // race this test targets — so all note creation happens first. + List noteIds = new ArrayList<>(); + for (int iter = 0; iter < 20; iter++) { + noteIds.add(newNoteWithParagraph("RaceNote" + iter, "race content " + iter)); + } + + for (int iter = 0; iter < 20; iter++) { + String noteId = noteIds.get(iter); + EmbeddingSearch search = openSearch(); + search.addNoteIndex(noteId); + + CountDownLatch startLatch = new CountDownLatch(2); + CountDownLatch doneLatch = new CountDownLatch(2); + Thread flushThread = new Thread(() -> { + awaitLatch(startLatch); + search.close(); // triggers flushIfDirty() -> saveNoteShard(noteId) synchronously + doneLatch.countDown(); + }); + Thread deleteThread = new Thread(() -> { + awaitLatch(startLatch); + search.deleteNoteIndex(noteId); + doneLatch.countDown(); + }); + flushThread.start(); + deleteThread.start(); + assertTrue(doneLatch.await(10, TimeUnit.SECONDS), + "iteration " + iter + ": both threads should finish"); + + // Whichever of the two ran first, settle any pending dirty mark the other one left + // behind (e.g. a delete that landed after the first flush already ran) the same + // deterministic way the rest of the suite does. + search.close(); + + EmbeddingSearch verify = openSearch(); // reload strictly from what's on disk + boolean resurrected = !verify.query("race content " + iter, id -> true).isEmpty(); + verify.close(); + assertFalse(resurrected, + "iteration " + iter + ": a concurrent flush must never recreate a deleted note's " + + "shard with stale, pre-delete content"); + } + } + + @Test + void failedMigrationDoesNotExposePartialShardSet() throws IOException { + File legacyFile = new File(indexDir, LEGACY_FILE); + String goodNoteId = "legacy-note-good"; + // 300 characters comfortably exceeds every common filesystem's per-component filename + // limit (255 bytes on ext4/APFS/etc.), so writing this note's shard fails with a real + // IOException partway through migration's staging loop — regardless of which of the + // two notes happens to be processed first (Set iteration order is unspecified), the + // whole migration must still fail cleanly rather than publish whatever got staged. + String badNoteId = "x".repeat(300); + writeLegacyIndexFile(legacyFile, + new LegacyEntry(goodNoteId + "/paragraph/p1", goodNoteId, "alpha content", "Alpha"), + new LegacyEntry(badNoteId + "/paragraph/p1", badNoteId, "beta content", "Beta")); + + EmbeddingSearch search = openSearch(); // migration attempt fails internally, swallowed + + File notesDir = new File(indexDir, NOTES_DIR); + File[] published = notesDir.exists() ? notesDir.listFiles() : null; + assertTrue(published == null || published.length == 0, + "a failed migration must never publish a partial shard set to notes/"); + search.close(); + } + + @Test + void migrationStopsWhenStaleStagingDirectoryCannotBeCleared() throws IOException { + File legacyFile = new File(indexDir, LEGACY_FILE); + writeLegacyIndexFile(legacyFile, + new LegacyEntry("legacy-note/paragraph/p1", "legacy-note", "legacy content", "Legacy")); + + Path stagingDir = indexDir.toPath().resolve("notes.migrating"); + Files.createDirectories(stagingDir); + Files.writeString(stagingDir.resolve("stale.bin"), "debris from an earlier migration"); + assumeTrue(Files.getFileStore(stagingDir).supportsFileAttributeView("posix")); + + Files.setPosixFilePermissions(stagingDir, PosixFilePermissions.fromString("r-x------")); + try { + assumeFalse(Files.isWritable(stagingDir), + "the test needs a staging directory whose children cannot be deleted"); + + EmbeddingSearch search = openSearch(); + + assertTrue(legacyFile.exists(), + "legacy index must remain when stale migration debris cannot be cleared"); + assertFalse(new File(indexDir, NOTES_DIR).exists(), + "uncleared staging contents must never be published as the live shard directory"); + search.close(); + } finally { + Files.setPosixFilePermissions(stagingDir, PosixFilePermissions.fromString("rwx------")); + } + } + + @Test + void legacyFileIsDeletedOnlyAfterCompleteMigration() throws IOException { + File legacyFile = new File(indexDir, LEGACY_FILE); + String noteId = newNoteWithParagraph("Delta", "current notebook content"); + + // Failing case: legacy file must survive so the next restart can retry. + String badNoteId = "y".repeat(300); + writeLegacyIndexFile(legacyFile, + new LegacyEntry(badNoteId + "/paragraph/p1", badNoteId, "gamma content", "Gamma")); + EmbeddingSearch failedAttempt = openSearch(); + assertTrue(legacyFile.exists(), "legacy file must not be deleted when migration fails"); + failedAttempt.close(); + + // Replace it with migratable data and confirm success deletes it. + assertTrue(legacyFile.delete()); + writeLegacyIndexFile(legacyFile, + new LegacyEntry(noteId + "/paragraph/p1", "Delta", "delta content", "Delta")); + EmbeddingSearch succeeded = openSearch(); + assertFalse(legacyFile.exists(), "legacy file must be deleted once migration fully succeeds"); + succeeded.close(); + } + + private static void awaitLatch(CountDownLatch latch) { + try { + latch.countDown(); + latch.await(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + + /** One entry for {@link #writeLegacyIndexFile}. */ + private static final class LegacyEntry { + final String docId; + final String noteName; + final String text; + final String title; + + LegacyEntry(String docId, String noteName, String text, String title) { + this.docId = docId; + this.noteName = noteName; + this.text = text; + this.title = title; + } + } + + /** Writes a shard whose header promises {@code 2} entries but only fully provides one — + * the first entry (containing {@code survivingText}) is written in full, then the + * second entry is cut off partway through its fields, so reading it throws. Reproduces + * "one or more entries read successfully, then a failure" rather than failing on the + * shard header itself. */ + private static void writeTruncatedShard(File file, String noteId, String survivingText) + throws IOException { + java.io.ByteArrayOutputStream baos = new java.io.ByteArrayOutputStream(); + try (DataOutputStream out = new DataOutputStream(baos)) { + out.writeInt(1); // SHARD_VERSION + out.writeUTF(noteId); + out.writeInt(2); // claims two entries + // Entry 1: fully valid. + out.writeUTF(noteId + "/paragraph/first-id"); + out.writeUTF("SomeNote"); + out.writeUTF(survivingText); + out.writeUTF(""); + out.writeUTF(""); + out.writeUTF(""); + for (int i = 0; i < 384; i++) { + out.writeFloat(0f); + } + // Entry 2: only the docId is written, then the stream ends — reading noteName next + // hits EOF partway through the entry. + out.writeUTF(noteId + "/paragraph/second-id"); + } + Files.write(file.toPath(), baos.toByteArray()); + } + + @Test + void corruptSingleShardRebuildsOnlyThatNoteOnNextStartup() + throws IOException, InterruptedException { + String goodNote = newNoteWithParagraph("GoodNote", "select * from good_table"); + String badNote = newNoteWithParagraph("BadNote", "select * from bad_table"); + + EmbeddingSearch search = openSearch(); + search.addNoteIndex(goodNote); + search.addNoteIndex(badNote); + search.close(); + + byte[] goodShardBefore = Files.readAllBytes(shardFile(goodNote).toPath()); + // Corrupt the bad note's shard: valid file, garbage version header. + Files.write(shardFile(badNote).toPath(), new byte[] {0, 0, 0, 99, 1, 2, 3}); + + EmbeddingSearch search2 = openSearch(); // constructor: loadAllShards() finds the corruption + assertFalse(shardFile(badNote).exists(), + "corrupt shard should be deleted rather than left on disk"); + + // Corrupt-shard rebuilds are scheduled via notebook.addInitConsumer, same mechanism as a + // normal server startup bootstrap — drive it the same way the server would. + notebook.initNotebook(); + assertTrue(notebook.waitForFinishInit(10, TimeUnit.SECONDS)); + + byte[] goodShardAfter = Files.readAllBytes(shardFile(goodNote).toPath()); + assertArrayEquals(goodShardBefore, goodShardAfter, + "the untouched note's shard must never be rewritten just because another note's " + + "shard was corrupt — that's the point of sharding by note"); + + List> results = search2.query("bad_table", id -> true); + assertTrue(results.stream().anyMatch(r -> r.get("text").contains("bad_table")), + "the corrupted note's content should have been rebuilt from the notebook"); + search2.close(); + } + + @Test + void missingSingleShardRebuildsOnlyThatNoteOnNextStartup() + throws IOException, InterruptedException { + String presentNote = newNoteWithParagraph("PresentNote", "select * from present_table"); + String missingNote = newNoteWithParagraph("MissingNote", "select * from missing_table"); + + EmbeddingSearch search = openSearch(); + search.addNoteIndex(presentNote); + search.addNoteIndex(missingNote); + search.close(); + + byte[] presentShardBefore = Files.readAllBytes(shardFile(presentNote).toPath()); + assertTrue(shardFile(missingNote).delete()); + + EmbeddingSearch reopened = openSearch(); + notebook.initNotebook(); + assertTrue(notebook.waitForFinishInit(10, TimeUnit.SECONDS)); + + assertArrayEquals(presentShardBefore, Files.readAllBytes(shardFile(presentNote).toPath()), + "an existing shard must remain unchanged when a different note's shard is missing"); + assertTrue(reopened.query("missing_table", id -> true).stream() + .anyMatch(r -> r.get("text").contains("missing_table")), + "the note whose shard vanished should be rebuilt from the notebook"); + + reopened.close(); + assertTrue(shardFile(missingNote).exists(), + "the rebuilt note should be persisted again on flush"); + } + + @Test + void orphanShardIsDeletedAndNotLoadedOnStartup() throws IOException { + String noteId = newNoteWithParagraph("RemovedNote", "content from a removed note"); + + EmbeddingSearch search = openSearch(); + search.addNoteIndex(noteId); + search.close(); + assertTrue(shardFile(noteId).exists()); + + // removeCorruptedNote deliberately has no NoteRemoveEvent. This also models a process + // stopping after notebook deletion was persisted but before the asynchronous search + // event or dirty-shard flush completed. + notebook.removeCorruptedNote(noteId, AuthenticationInfo.ANONYMOUS); + + EmbeddingSearch reopened = openSearch(); + + assertFalse(shardFile(noteId).exists(), + "a shard whose note no longer exists should be removed during startup"); + assertTrue(reopened.query("content from a removed note", id -> true).isEmpty(), + "orphan shard entries must not be loaded into the in-memory search index"); + reopened.close(); + } + + @Test + void concurrentParagraphMutationsOnSameNoteStayConsistent() throws Exception { + String noteId = notebook.createNote("Concurrent", AuthenticationInfo.ANONYMOUS); + int paragraphCount = 20; + List paragraphIds = new ArrayList<>(); + notebook.processNote(noteId, note -> { + for (int i = 0; i < paragraphCount; i++) { + Paragraph p = note.addNewParagraph(AuthenticationInfo.ANONYMOUS); + p.setText("paragraph " + i); + paragraphIds.add(p.getId()); + } + return null; + }); + + EmbeddingSearch search = openSearch(); + search.addNoteIndex(noteId); + + // Half the paragraphs get deleted, half get re-indexed, all concurrently on one note. + // deleteParagraphIndex previously took no lock at all — this is the regression test for + // that fix. + ExecutorService pool = Executors.newFixedThreadPool(8); + CountDownLatch done = new CountDownLatch(paragraphCount); + for (int i = 0; i < paragraphCount; i++) { + final String pid = paragraphIds.get(i); + final boolean delete = i % 2 == 0; + pool.submit(() -> { + try { + if (delete) { + search.deleteParagraphIndex(noteId, pid); + } else { + search.updateParagraphIndex(noteId, pid); + } + } finally { + done.countDown(); + } + }); + } + assertTrue(done.await(30, TimeUnit.SECONDS), "all concurrent mutations should complete"); + pool.shutdown(); + + search.close(); // must not throw, and must persist a consistent shard + + EmbeddingSearch search2 = openSearch(); + int remaining = 0; + for (int i = 1; i < paragraphCount; i += 2) { + final String expectedText = "paragraph " + i; + List> r = search2.query(expectedText, id -> true); + if (r.stream().anyMatch(res -> res.get("text").contains(expectedText))) { + remaining++; + } + } + assertEquals(paragraphCount / 2, remaining, + "every non-deleted paragraph should have survived the concurrent mutations"); + search2.close(); + } + + @Test + void flushOnlyWritesDirtyNotes() throws IOException { + String noteA = newNoteWithParagraph("NoteA", "select * from a"); + String noteB = newNoteWithParagraph("NoteB", "select * from b"); + + EmbeddingSearch search = openSearch(); + search.addNoteIndex(noteA); + search.addNoteIndex(noteB); + search.close(); + + byte[] noteBBefore = Files.readAllBytes(shardFile(noteB).toPath()); + + EmbeddingSearch search2 = openSearch(); + String paragraphId = lastParagraphId(noteA); + notebook.processNote(noteA, note -> { + note.getLastParagraph().setText("changed only in note A"); + return null; + }); + search2.updateParagraphIndex(noteA, paragraphId); + // Note B was never touched, so it was never added to dirtyNoteIds. + search2.close(); + + byte[] noteBAfter = Files.readAllBytes(shardFile(noteB).toPath()); + assertArrayEquals(noteBBefore, noteBAfter, "non-dirty note must not be rewritten on flush"); + } + + /** Writes a legacy (pre-ZEPPELIN-6412) single-file index with one or more entries, + * matching the format {@code EmbeddingSearch.loadLegacyIndex} reads: + * [int version=3][int count] then per entry [utf docId][utf noteName][utf text] + * [utf title][utf tables][utf output][float[384] embedding]. */ + private static void writeLegacyIndexFile(File file, LegacyEntry... entries) throws IOException { + try (DataOutputStream out = new DataOutputStream(Files.newOutputStream(file.toPath()))) { + out.writeInt(3); // legacy INDEX_VERSION + out.writeInt(entries.length); + for (LegacyEntry e : entries) { + out.writeUTF(e.docId); + out.writeUTF(e.noteName); + out.writeUTF(e.text); + out.writeUTF(e.title); + out.writeUTF(""); // tables + out.writeUTF(""); // output + for (int i = 0; i < 384; i++) { + out.writeFloat(0f); + } + } + } + } +}