diff --git a/.gitignore b/.gitignore index 4662ff4c5a4d1..f5637f04366fc 100644 --- a/.gitignore +++ b/.gitignore @@ -70,3 +70,4 @@ testfixtures_shared/ # build files generated doc-tools/missing-doclet/bin/ /sandbox/plugins/engine-datafusion/target/ +**/Cargo.lock \ No newline at end of file diff --git a/CHANGELOG.md b/CHANGELOG.md index 356bc0f8cd93d..fe0ecf102f6b2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -27,6 +27,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), - Add support for enabling pluggable data formats, starting with phase-1 of decoupling shard from engine, and introducing basic abstractions ([#20675](https://github.com/opensearch-project/OpenSearch/pull/20675)) - Add concurrent queue in libs and composite engine sandbox plugin ([#20909](https://github.com/opensearch-project/OpenSearch/pull/20909)) - Add interface for the Multi format merge flow ([#20908](https://github.com/opensearch-project/OpenSearch/pull/20908)) +- Add CatalogSnapshotManager lifecycle management with reference-counted snapshot tracking and serialization support for Segment and WriterFileSet ([#20982](https://github.com/opensearch-project/OpenSearch/pull/20982)) - Add warmup phase to wait for lag to catch up in pull-based ingestion before serving ([#20526](https://github.com/opensearch-project/OpenSearch/pull/20526)) - Add a new static method to IndicesOptions API to expose `STRICT_EXPAND_OPEN_HIDDEN_FORBID_CLOSED` index option ([#20980](https://github.com/opensearch-project/OpenSearch/pull/20980)) diff --git a/sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/DatafusionReaderManager.java b/sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/DatafusionReaderManager.java index f97f11f78b743..591ca9f26135e 100644 --- a/sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/DatafusionReaderManager.java +++ b/sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/DatafusionReaderManager.java @@ -10,8 +10,8 @@ import org.opensearch.common.annotation.ExperimentalApi; import org.opensearch.index.engine.dataformat.DataFormat; -import org.opensearch.index.engine.exec.CatalogSnapshot; import org.opensearch.index.engine.exec.EngineReaderManager; +import org.opensearch.index.engine.exec.coord.CatalogSnapshot; import org.opensearch.index.shard.ShardPath; import java.io.IOException; diff --git a/sandbox/plugins/analytics-backend-lucene/src/main/java/org/opensearch/be/lucene/LuceneReaderManager.java b/sandbox/plugins/analytics-backend-lucene/src/main/java/org/opensearch/be/lucene/LuceneReaderManager.java index c46d480bccfb3..d3fe2c6338089 100644 --- a/sandbox/plugins/analytics-backend-lucene/src/main/java/org/opensearch/be/lucene/LuceneReaderManager.java +++ b/sandbox/plugins/analytics-backend-lucene/src/main/java/org/opensearch/be/lucene/LuceneReaderManager.java @@ -12,8 +12,8 @@ import org.apache.lucene.search.ReferenceManager; import org.opensearch.common.annotation.ExperimentalApi; import org.opensearch.index.engine.dataformat.DataFormat; -import org.opensearch.index.engine.exec.CatalogSnapshot; import org.opensearch.index.engine.exec.EngineReaderManager; +import org.opensearch.index.engine.exec.coord.CatalogSnapshot; import java.io.IOException; import java.util.Collection; diff --git a/sandbox/plugins/analytics-engine/src/test/java/org/opensearch/analytics/exec/DefaultPlanExecutorTests.java b/sandbox/plugins/analytics-engine/src/test/java/org/opensearch/analytics/exec/DefaultPlanExecutorTests.java index d623993227d95..b7c832438092f 100644 --- a/sandbox/plugins/analytics-engine/src/test/java/org/opensearch/analytics/exec/DefaultPlanExecutorTests.java +++ b/sandbox/plugins/analytics-engine/src/test/java/org/opensearch/analytics/exec/DefaultPlanExecutorTests.java @@ -30,15 +30,17 @@ import org.opensearch.cluster.metadata.IndexMetadata; import org.opensearch.cluster.metadata.Metadata; import org.opensearch.cluster.service.ClusterService; +import org.opensearch.common.concurrent.GatedCloseable; import org.opensearch.core.index.Index; import org.opensearch.index.IndexService; import org.opensearch.index.engine.DataFormatAwareEngine; import org.opensearch.index.engine.dataformat.DataFormat; import org.opensearch.index.engine.dataformat.FieldTypeCapabilities; -import org.opensearch.index.engine.exec.CatalogSnapshot; import org.opensearch.index.engine.exec.EngineReaderManager; import org.opensearch.index.engine.exec.Segment; import org.opensearch.index.engine.exec.WriterFileSet; +import org.opensearch.index.engine.exec.coord.CatalogSnapshot; +import org.opensearch.index.engine.exec.coord.CatalogSnapshotManager; import org.opensearch.index.shard.IndexShard; import org.opensearch.indices.IndicesService; import org.opensearch.test.OpenSearchTestCase; @@ -104,13 +106,15 @@ public void testEndToEndExecuteWithMockBackend() throws IOException { Segment seg1 = Segment.builder(0L).addSearchableFiles(format, fs1).build(); Segment seg2 = Segment.builder(1L).addSearchableFiles(format, fs2).build(); - MockCatalogSnapshot snapshot = new MockCatalogSnapshot(1L, List.of(seg1, seg2), format); + + CatalogSnapshotManager snapshotManager = new CatalogSnapshotManager(1L, 1L, 0L, List.of(seg1, seg2), 2L, Map.of()); MockReaderManager readerManager = new MockReaderManager(format.name()); - readerManager.afterRefresh(true, snapshot); + try (GatedCloseable ref = snapshotManager.acquireSnapshot()) { + readerManager.afterRefresh(true, ref.get()); + } - DataFormatAwareEngine engine = new DataFormatAwareEngine(Map.of(format, readerManager)); - engine.setLatestSnapshot(snapshot); + DataFormatAwareEngine engine = new DataFormatAwareEngine(Map.of(format, readerManager), snapshotManager); // Mock shard + cluster wiring IndexShard shard = mock(IndexShard.class); @@ -288,16 +292,18 @@ public String serializeToString() { } @Override - public void setCatalogSnapshotMap(Map map) {} - - @Override - public void setUserData(Map userData, boolean b) {} + public void setUserData(Map userData) {} @Override public Object getReader(DataFormat dataFormat) { return null; } + @Override + public MockCatalogSnapshot clone() { + return new MockCatalogSnapshot(generation, segments, format); + } + @Override protected void closeInternal() {} } diff --git a/server/src/main/java/org/opensearch/index/engine/DataFormatAwareEngine.java b/server/src/main/java/org/opensearch/index/engine/DataFormatAwareEngine.java index 1d55e84d68b37..03d4354f9d03d 100644 --- a/server/src/main/java/org/opensearch/index/engine/DataFormatAwareEngine.java +++ b/server/src/main/java/org/opensearch/index/engine/DataFormatAwareEngine.java @@ -11,10 +11,11 @@ import org.opensearch.common.annotation.ExperimentalApi; import org.opensearch.common.concurrent.GatedCloseable; import org.opensearch.index.engine.dataformat.DataFormat; -import org.opensearch.index.engine.exec.CatalogSnapshot; import org.opensearch.index.engine.exec.DataFormatAwareEngineFactory; import org.opensearch.index.engine.exec.EngineReaderManager; import org.opensearch.index.engine.exec.IndexReaderProvider; +import org.opensearch.index.engine.exec.coord.CatalogSnapshot; +import org.opensearch.index.engine.exec.coord.CatalogSnapshotManager; import java.io.Closeable; import java.io.IOException; @@ -35,51 +36,45 @@ public class DataFormatAwareEngine implements IndexReaderProvider, Closeable { private final Map> readerManagers; - private volatile CatalogSnapshot latestSnapshot; + private volatile CatalogSnapshotManager catalogSnapshotManager; /** - * Constructs a new DataFormatAwareEngine with pre-built maps. + * Constructs a new DataFormatAwareEngine. * Prefer using {@link DataFormatAwareEngineFactory#create()}. */ + public DataFormatAwareEngine(Map> readerManagers, CatalogSnapshotManager catalogSnapshotManager) { + this.readerManagers = readerManagers; + this.catalogSnapshotManager = catalogSnapshotManager; + } + + /** + * Constructs a new DataFormatAwareEngine without a snapshot manager. + * The manager must be set via {@link #setCatalogSnapshotManager} before acquiring readers. + */ public DataFormatAwareEngine(Map> readerManagers) { this.readerManagers = readerManagers; } - public EngineReaderManager getReaderManager(DataFormat format) { - return readerManagers.get(format); + public void setCatalogSnapshotManager(CatalogSnapshotManager catalogSnapshotManager) { + this.catalogSnapshotManager = catalogSnapshotManager; } - /** - * Called by the catalog snapshot lifecycle listener after a refresh - * to update the latest searchable snapshot. - */ - public void setLatestSnapshot(CatalogSnapshot snapshot) { - CatalogSnapshot prev = this.latestSnapshot; - this.latestSnapshot = snapshot; - if (prev != null) { - prev.decRef(); - } + public EngineReaderManager getReaderManager(DataFormat format) { + return readerManagers.get(format); } /** * Acquires a DataFormatAwareReader on the latest catalog snapshot. - * The snapshot is incRef'd; the caller MUST close the returned - * {@link DataFormatAwareReader} when done, which decRef's the snapshot. + * The caller MUST close the returned {@link DataFormatAwareReader} when done, + * which releases the snapshot reference. */ public GatedCloseable acquireReader() throws IOException { - CatalogSnapshot snapshot = latestSnapshot; - if (snapshot == null) { - throw new IllegalStateException("No catalog snapshot available"); + if (catalogSnapshotManager == null) { + throw new IllegalStateException("CatalogSnapshotManager not set"); } - return acquireReader(snapshot); - } - - /** - * Acquires a dataFormatAwareReader on a specific catalog snapshot. - */ - private GatedCloseable acquireReader(CatalogSnapshot catalogSnapshot) throws IOException { - catalogSnapshot.incRef(); + GatedCloseable snapshotRef = catalogSnapshotManager.acquireSnapshot(); try { + CatalogSnapshot catalogSnapshot = snapshotRef.get(); Map readers = new HashMap<>(); for (Map.Entry> entry : readerManagers.entrySet()) { Object reader = entry.getValue().getReader(catalogSnapshot); @@ -87,10 +82,10 @@ private GatedCloseable acquireReader(CatalogSnapshot catalogSnapshot) th readers.put(entry.getKey(), reader); } } - DataFormatAwareReader reader = new DataFormatAwareReader(catalogSnapshot, readers); + DataFormatAwareReader reader = new DataFormatAwareReader(catalogSnapshot, snapshotRef, readers); return new GatedCloseable<>(reader, reader::close); } catch (Exception e) { - catalogSnapshot.decRef(); + snapshotRef.close(); throw e; } } @@ -102,10 +97,16 @@ private GatedCloseable acquireReader(CatalogSnapshot catalogSnapshot) th @ExperimentalApi public static class DataFormatAwareReader implements IndexReaderProvider.Reader { private final CatalogSnapshot catalogSnapshot; + private final GatedCloseable snapshotRef; private final Map readers; - DataFormatAwareReader(CatalogSnapshot catalogSnapshot, Map readers) { + DataFormatAwareReader( + CatalogSnapshot catalogSnapshot, + GatedCloseable snapshotRef, + Map readers + ) { this.catalogSnapshot = catalogSnapshot; + this.snapshotRef = snapshotRef; this.readers = readers; } @@ -121,7 +122,11 @@ public CatalogSnapshot catalogSnapshot() { @Override public void close() { - catalogSnapshot.decRef(); + try { + snapshotRef.close(); + } catch (IOException e) { + throw new RuntimeException("Failed to release catalog snapshot reference", e); + } } } diff --git a/server/src/main/java/org/opensearch/index/engine/EngineBackedIndexer.java b/server/src/main/java/org/opensearch/index/engine/EngineBackedIndexer.java index 4fd056d97c762..eb03ff54c0e11 100644 --- a/server/src/main/java/org/opensearch/index/engine/EngineBackedIndexer.java +++ b/server/src/main/java/org/opensearch/index/engine/EngineBackedIndexer.java @@ -14,8 +14,8 @@ import org.opensearch.common.unit.TimeValue; import org.opensearch.core.common.unit.ByteSizeValue; import org.opensearch.index.VersionType; -import org.opensearch.index.engine.exec.CatalogSnapshot; import org.opensearch.index.engine.exec.Indexer; +import org.opensearch.index.engine.exec.coord.CatalogSnapshot; import org.opensearch.index.mapper.DocumentMapperForType; import org.opensearch.index.mapper.SourceToParse; import org.opensearch.index.merge.MergeStats; diff --git a/server/src/main/java/org/opensearch/index/engine/dataformat/merge/MergeHandler.java b/server/src/main/java/org/opensearch/index/engine/dataformat/merge/MergeHandler.java index 48dfe22952202..7c6b2e3cb657d 100644 --- a/server/src/main/java/org/opensearch/index/engine/dataformat/merge/MergeHandler.java +++ b/server/src/main/java/org/opensearch/index/engine/dataformat/merge/MergeHandler.java @@ -15,9 +15,9 @@ import org.opensearch.common.logging.Loggers; import org.opensearch.core.index.shard.ShardId; import org.opensearch.index.engine.dataformat.MergeResult; -import org.opensearch.index.engine.exec.CatalogSnapshot; import org.opensearch.index.engine.exec.Indexer; import org.opensearch.index.engine.exec.Segment; +import org.opensearch.index.engine.exec.coord.CatalogSnapshot; import java.util.ArrayDeque; import java.util.Collection; diff --git a/server/src/main/java/org/opensearch/index/engine/exec/CatalogSnapshotLifecycleListener.java b/server/src/main/java/org/opensearch/index/engine/exec/CatalogSnapshotLifecycleListener.java index e0a40709acf33..ee0431c3d2a63 100644 --- a/server/src/main/java/org/opensearch/index/engine/exec/CatalogSnapshotLifecycleListener.java +++ b/server/src/main/java/org/opensearch/index/engine/exec/CatalogSnapshotLifecycleListener.java @@ -9,6 +9,7 @@ package org.opensearch.index.engine.exec; import org.opensearch.common.annotation.ExperimentalApi; +import org.opensearch.index.engine.exec.coord.CatalogSnapshot; import java.io.IOException; diff --git a/server/src/main/java/org/opensearch/index/engine/exec/DataFormatEngineCatalogSnapshotListener.java b/server/src/main/java/org/opensearch/index/engine/exec/DataFormatEngineCatalogSnapshotListener.java index 85e247bd29fd1..c8c6a2dc2f002 100644 --- a/server/src/main/java/org/opensearch/index/engine/exec/DataFormatEngineCatalogSnapshotListener.java +++ b/server/src/main/java/org/opensearch/index/engine/exec/DataFormatEngineCatalogSnapshotListener.java @@ -10,6 +10,7 @@ import org.opensearch.common.annotation.ExperimentalApi; import org.opensearch.index.engine.dataformat.DataFormat; +import org.opensearch.index.engine.exec.coord.CatalogSnapshot; import java.io.IOException; import java.util.Collection; diff --git a/server/src/main/java/org/opensearch/index/engine/exec/EngineReaderManager.java b/server/src/main/java/org/opensearch/index/engine/exec/EngineReaderManager.java index b420dd6299471..b7d412850bd3d 100644 --- a/server/src/main/java/org/opensearch/index/engine/exec/EngineReaderManager.java +++ b/server/src/main/java/org/opensearch/index/engine/exec/EngineReaderManager.java @@ -9,6 +9,7 @@ package org.opensearch.index.engine.exec; import org.opensearch.common.annotation.ExperimentalApi; +import org.opensearch.index.engine.exec.coord.CatalogSnapshot; import java.io.IOException; diff --git a/server/src/main/java/org/opensearch/index/engine/exec/IndexFileDeleter.java b/server/src/main/java/org/opensearch/index/engine/exec/IndexFileDeleter.java index 61507b7ffe9d7..214402e4064b5 100644 --- a/server/src/main/java/org/opensearch/index/engine/exec/IndexFileDeleter.java +++ b/server/src/main/java/org/opensearch/index/engine/exec/IndexFileDeleter.java @@ -11,6 +11,7 @@ import org.opensearch.common.annotation.ExperimentalApi; import org.opensearch.index.engine.DataFormatAwareEngine; import org.opensearch.index.engine.dataformat.DataFormat; +import org.opensearch.index.engine.exec.coord.CatalogSnapshot; import org.opensearch.index.shard.ShardPath; import java.io.IOException; diff --git a/server/src/main/java/org/opensearch/index/engine/exec/IndexReaderProvider.java b/server/src/main/java/org/opensearch/index/engine/exec/IndexReaderProvider.java index 5ecd9317f40a9..0662fc40c4833 100644 --- a/server/src/main/java/org/opensearch/index/engine/exec/IndexReaderProvider.java +++ b/server/src/main/java/org/opensearch/index/engine/exec/IndexReaderProvider.java @@ -11,6 +11,7 @@ import org.opensearch.common.annotation.ExperimentalApi; import org.opensearch.common.concurrent.GatedCloseable; import org.opensearch.index.engine.dataformat.DataFormat; +import org.opensearch.index.engine.exec.coord.CatalogSnapshot; import java.io.Closeable; import java.io.IOException; diff --git a/server/src/main/java/org/opensearch/index/engine/exec/Indexer.java b/server/src/main/java/org/opensearch/index/engine/exec/Indexer.java index 39b6ea4e80a88..d2025a2f0c22e 100644 --- a/server/src/main/java/org/opensearch/index/engine/exec/Indexer.java +++ b/server/src/main/java/org/opensearch/index/engine/exec/Indexer.java @@ -16,6 +16,7 @@ import org.opensearch.index.engine.EngineException; import org.opensearch.index.engine.LifecycleAware; import org.opensearch.index.engine.SafeCommitInfo; +import org.opensearch.index.engine.exec.coord.CatalogSnapshot; import org.opensearch.index.translog.Translog; import org.opensearch.index.translog.TranslogManager; diff --git a/server/src/main/java/org/opensearch/index/engine/exec/IndexerLifecycleOperations.java b/server/src/main/java/org/opensearch/index/engine/exec/IndexerLifecycleOperations.java index e26189444d33a..6d51587e49def 100644 --- a/server/src/main/java/org/opensearch/index/engine/exec/IndexerLifecycleOperations.java +++ b/server/src/main/java/org/opensearch/index/engine/exec/IndexerLifecycleOperations.java @@ -12,6 +12,7 @@ import org.opensearch.common.unit.TimeValue; import org.opensearch.core.common.unit.ByteSizeValue; import org.opensearch.index.engine.EngineException; +import org.opensearch.index.engine.exec.coord.CatalogSnapshot; import org.opensearch.index.shard.ShardPath; import java.io.IOException; diff --git a/server/src/main/java/org/opensearch/index/engine/exec/Segment.java b/server/src/main/java/org/opensearch/index/engine/exec/Segment.java index 5a5811804aff3..576d871832dde 100644 --- a/server/src/main/java/org/opensearch/index/engine/exec/Segment.java +++ b/server/src/main/java/org/opensearch/index/engine/exec/Segment.java @@ -9,23 +9,58 @@ package org.opensearch.index.engine.exec; import org.opensearch.common.annotation.ExperimentalApi; +import org.opensearch.core.common.io.stream.StreamInput; +import org.opensearch.core.common.io.stream.StreamOutput; +import org.opensearch.core.common.io.stream.Writeable; import org.opensearch.index.engine.dataformat.DataFormat; +import java.io.IOException; import java.util.HashMap; import java.util.Map; +import java.util.function.Function; /** * Represents a segment in the catalog snapshot containing files grouped by data format. * Each segment has a unique generation number and maintains searchable files organized by their data format type. - * This class is serializable and can be transmitted across nodes for replication and recovery operations. */ @ExperimentalApi -public record Segment(long generation, Map dfGroupedSearchableFiles) { +public record Segment(long generation, Map dfGroupedSearchableFiles) implements Writeable { public Segment { dfGroupedSearchableFiles = Map.copyOf(dfGroupedSearchableFiles); } + /** + * Constructs a Segment by deserializing from a {@link StreamInput}. + * + * @param in the stream input to read from + * @param directoryResolver function that maps a data format name to its directory path + */ + public Segment(StreamInput in, Function directoryResolver) throws IOException { + this(in.readLong(), readWriterFileSets(in, directoryResolver)); + } + + private static Map readWriterFileSets(StreamInput in, Function directoryResolver) + throws IOException { + int size = in.readVInt(); + Map map = new HashMap<>(size); + for (int i = 0; i < size; i++) { + String key = in.readString(); + map.put(key, new WriterFileSet(in, directoryResolver.apply(key))); + } + return map; + } + + @Override + public void writeTo(StreamOutput out) throws IOException { + out.writeLong(generation); + out.writeVInt(dfGroupedSearchableFiles.size()); + for (Map.Entry entry : dfGroupedSearchableFiles.entrySet()) { + out.writeString(entry.getKey()); + entry.getValue().writeTo(out); + } + } + public static Builder builder(long generation) { return new Builder(generation); } diff --git a/server/src/main/java/org/opensearch/index/engine/exec/WriterFileSet.java b/server/src/main/java/org/opensearch/index/engine/exec/WriterFileSet.java index 65fe34410e9d0..5c497d21cfc74 100644 --- a/server/src/main/java/org/opensearch/index/engine/exec/WriterFileSet.java +++ b/server/src/main/java/org/opensearch/index/engine/exec/WriterFileSet.java @@ -9,6 +9,9 @@ package org.opensearch.index.engine.exec; import org.opensearch.common.annotation.ExperimentalApi; +import org.opensearch.core.common.io.stream.StreamInput; +import org.opensearch.core.common.io.stream.StreamOutput; +import org.opensearch.core.common.io.stream.Writeable; import java.io.IOException; import java.nio.file.Path; @@ -20,12 +23,19 @@ * Groups files by directory and writer generation, tracking metadata such as row count and total size. */ @ExperimentalApi -public record WriterFileSet(String directory, long writerGeneration, Set files, long numRows) { +public record WriterFileSet(String directory, long writerGeneration, Set files, long numRows) implements Writeable { public WriterFileSet { files = Set.copyOf(files); } + /** + * Constructs a WriterFileSet by deserializing from a {@link StreamInput}. + */ + public WriterFileSet(StreamInput in, String directory) throws IOException { + this(directory, in.readLong(), new HashSet<>(in.readStringList()), in.readLong()); + } + public long getTotalSize() { return files.stream().mapToLong(file -> { try { @@ -41,6 +51,15 @@ public String toString() { return "WriterFileSet{" + "directory=" + directory + ", writerGeneration=" + writerGeneration + ", files=" + files + '}'; } + /** + * Serializes this WriterFileSet to the given stream output. + */ + public void writeTo(StreamOutput out) throws IOException { + out.writeLong(writerGeneration); + out.writeStringCollection(files); + out.writeLong(numRows); + } + /** * Creates a new builder for constructing WriterFileSet instances. * diff --git a/server/src/main/java/org/opensearch/index/engine/exec/CatalogSnapshot.java b/server/src/main/java/org/opensearch/index/engine/exec/coord/CatalogSnapshot.java similarity index 52% rename from server/src/main/java/org/opensearch/index/engine/exec/CatalogSnapshot.java rename to server/src/main/java/org/opensearch/index/engine/exec/coord/CatalogSnapshot.java index 80abcb59eccbe..f10cd55a075e3 100644 --- a/server/src/main/java/org/opensearch/index/engine/exec/CatalogSnapshot.java +++ b/server/src/main/java/org/opensearch/index/engine/exec/coord/CatalogSnapshot.java @@ -6,11 +6,16 @@ * compatible open source license. */ -package org.opensearch.index.engine.exec; +package org.opensearch.index.engine.exec.coord; import org.opensearch.common.annotation.ExperimentalApi; import org.opensearch.common.util.concurrent.AbstractRefCounted; +import org.opensearch.core.common.io.stream.StreamInput; +import org.opensearch.core.common.io.stream.StreamOutput; +import org.opensearch.core.common.io.stream.Writeable; import org.opensearch.index.engine.dataformat.DataFormat; +import org.opensearch.index.engine.exec.Segment; +import org.opensearch.index.engine.exec.WriterFileSet; import java.io.IOException; import java.util.Collection; @@ -21,11 +26,17 @@ /** * Abstract base class representing a snapshot of the catalog state at a specific point in time. * Maintains versioned information about segments, files, and metadata for index operations. - * Extends AbstractRefCounted to support reference counting for safe concurrent access. - * Subclasses must implement methods for accessing file metadata, segments, and user data. + * Uses an internal reference counter for safe concurrent access. + * Subclasses must implement {@link #closeInternal()} for resource cleanup and methods for + * accessing file metadata, segments, and user data. + * + *

Important: Do not call {@code incRef()}, {@code decRef()}, or {@code tryIncRef()} directly. + * Use {@link org.opensearch.index.engine.exec.coord.CatalogSnapshotManager#acquireSnapshot()} to obtain + * a reference-counted handle, and close the returned {@link org.opensearch.common.concurrent.GatedCloseable} + * when done. The manager handles all reference counting internally.

*/ @ExperimentalApi -public abstract class CatalogSnapshot extends AbstractRefCounted { +public abstract class CatalogSnapshot implements Writeable, Cloneable { /** * Key for storing catalog snapshot in user data. @@ -43,10 +54,40 @@ public abstract class CatalogSnapshot extends AbstractRefCounted { protected final long generation; protected long version; - public CatalogSnapshot(String name, long generation, long version) { - super(name); + private final AbstractRefCounted refCounter; + + protected CatalogSnapshot(String name, long generation, long version) { this.generation = generation; this.version = version; + this.refCounter = new AbstractRefCounted(name) { + @Override + protected void closeInternal() { + CatalogSnapshot.this.closeInternal(); + } + }; + } + + /** + * Constructs a CatalogSnapshot from a {@link StreamInput}. + * + * @param in the stream input to read from + * @throws IOException if an I/O error occurs + */ + protected CatalogSnapshot(StreamInput in) throws IOException { + this.generation = in.readLong(); + this.version = in.readLong(); + this.refCounter = new AbstractRefCounted("catalog_snapshot") { + @Override + protected void closeInternal() { + CatalogSnapshot.this.closeInternal(); + } + }; + } + + @Override + public void writeTo(StreamOutput out) throws IOException { + out.writeLong(generation); + out.writeLong(version); } public long getGeneration() { @@ -57,6 +98,35 @@ public long getVersion() { return version; } + // Package-private ref counting — only accessible within exec.coord (i.e., CatalogSnapshotManager) + + /** + * Decrements the reference count. Returns {@code true} if the count reached zero + * and {@link #closeInternal()} was invoked. + */ + boolean decRef() { + return refCounter.decRef(); + } + + /** + * Tries to increment the reference count. Returns {@code false} if the snapshot is already closed. + */ + boolean tryIncRef() { + return refCounter.tryIncRef(); + } + + /** + * Returns the current reference count. + */ + int refCount() { + return refCounter.refCount(); + } + + /** + * Called when the reference count reaches zero. Subclasses should release any resources here. + */ + protected abstract void closeInternal(); + /** * Gets user-defined metadata associated with this catalog snapshot. * @@ -108,13 +178,6 @@ public long getVersion() { */ public abstract String serializeToString() throws IOException; - /** - * Sets the catalog snapshot map for tracking multiple snapshots. - * - * @param catalogSnapshotMap map of generation to catalog snapshots - */ - public abstract void setCatalogSnapshotMap(Map catalogSnapshotMap); - /** * Creates a clone without acquiring a reference count. * Used for Lucene compatibility where clone is required. @@ -122,8 +185,6 @@ public long getVersion() { * @return this catalog snapshot instance */ public CatalogSnapshot cloneNoAcquire() { - // Still using the clone call since Lucene call requires clone. This will allow a SegmentsInfos backed CatalogSnapshot to use the - // same method in calls. return this; } @@ -131,9 +192,16 @@ public CatalogSnapshot cloneNoAcquire() { * Sets user-defined metadata for this catalog snapshot. * * @param userData map of user data key-value pairs - * @param b additional boolean parameter for implementation-specific behavior */ - public abstract void setUserData(Map userData, boolean b); + public abstract void setUserData(Map userData); + + /** + * Creates a deep copy of this catalog snapshot. The cloned snapshot starts with a fresh reference count of 1. + * Subclasses must ensure all mutable state is properly copied. + * + * @return a new {@link CatalogSnapshot} with the same logical state + */ + public abstract CatalogSnapshot clone(); public abstract Object getReader(DataFormat dataFormat); } diff --git a/server/src/main/java/org/opensearch/index/engine/exec/coord/CatalogSnapshotManager.java b/server/src/main/java/org/opensearch/index/engine/exec/coord/CatalogSnapshotManager.java new file mode 100644 index 0000000000000..1005c5e2b9eb6 --- /dev/null +++ b/server/src/main/java/org/opensearch/index/engine/exec/coord/CatalogSnapshotManager.java @@ -0,0 +1,131 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +package org.opensearch.index.engine.exec.coord; + +import org.opensearch.common.annotation.ExperimentalApi; +import org.opensearch.common.concurrent.GatedCloseable; +import org.opensearch.index.engine.exec.Segment; + +import java.io.Closeable; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.atomic.AtomicBoolean; + +/** + * Manages the lifecycle of {@link CatalogSnapshot} instances for the composite multi-format engine. + * + *

Tracks all live snapshots in a map keyed by generation. When a snapshot's reference count reaches + * zero (via {@link #decRefAndRemove}), it is automatically removed from the map. All {@code decRef} + * calls on managed snapshots go through this method to ensure consistent cleanup.

+ * + *

The write path (commit) is single-threaded (refresh is serialized per shard), while the read + * path (acquireSnapshot) is safe for concurrent access via volatile reads and {@code tryIncRef}.

+ */ +@ExperimentalApi +public class CatalogSnapshotManager implements Closeable { + + private volatile CatalogSnapshot latestCatalogSnapshot; + private final AtomicBoolean closed = new AtomicBoolean(false); + private final Map catalogSnapshotMap = new ConcurrentHashMap<>(); + + /** + * Constructs a new CatalogSnapshotManager with an initial snapshot built from the given parameters. + * + * @param id the unique snapshot identifier + * @param generation the initial generation number + * @param version the schema version + * @param segments the initial segments + * @param lastWriterGeneration the last writer generation + * @param userData user-defined metadata + */ + public CatalogSnapshotManager( + long id, + long generation, + long version, + List segments, + long lastWriterGeneration, + Map userData + ) { + DataformatAwareCatalogSnapshot initialSnapshot = new DataformatAwareCatalogSnapshot( + id, + generation, + version, + segments, + lastWriterGeneration, + userData + ); + this.latestCatalogSnapshot = initialSnapshot; + catalogSnapshotMap.put(initialSnapshot.getGeneration(), initialSnapshot); + } + + /** + * Acquires the current snapshot with an incremented reference count, wrapped in a {@link GatedCloseable} + * that calls {@link #decRefAndRemove} on close. + * + * @return a {@link GatedCloseable} wrapping the current {@link CatalogSnapshot} + * @throws IllegalStateException if the manager or snapshot is already closed + */ + public GatedCloseable acquireSnapshot() { + if (closed.get()) { + throw new IllegalStateException("CatalogSnapshotManager is closed"); + } + final CatalogSnapshot snapshot = latestCatalogSnapshot; + if (snapshot.tryIncRef() == false) { + throw new IllegalStateException("CatalogSnapshot [gen=" + snapshot.getGeneration() + "] is already closed"); + } + return new GatedCloseable<>(snapshot, () -> decRefAndRemove(snapshot)); + } + + /** + * Commits a new snapshot built from the given refreshed segments, replacing the current one. + * The new snapshot inherits user data from the current snapshot and increments the generation. + * The old snapshot is decRef'd and removed from the map if its count reaches zero. + * + * @param refreshedSegments the segments produced by the latest refresh + */ + public synchronized void commitNewSnapshot(List refreshedSegments) { + assert closed.get() == false : "Cannot commit to a closed CatalogSnapshotManager"; + + DataformatAwareCatalogSnapshot newSnapshot = new DataformatAwareCatalogSnapshot( + latestCatalogSnapshot.getId() + 1, + latestCatalogSnapshot.getGeneration() + 1, + latestCatalogSnapshot.getVersion(), + refreshedSegments, + latestCatalogSnapshot.getLastWriterGeneration() + 1, + latestCatalogSnapshot.getUserData() + ); + + CatalogSnapshot oldSnapshot = latestCatalogSnapshot; + latestCatalogSnapshot = newSnapshot; + decRefAndRemove(oldSnapshot); + } + + /** + * Decrements the reference count and removes the snapshot from the tracking map if it reaches zero. + * Generation is captured before decRef to avoid accessing the snapshot after closeInternal. + */ + private void decRefAndRemove(CatalogSnapshot snapshot) { + final long gen = snapshot.getGeneration(); + if (snapshot.decRef()) { + catalogSnapshotMap.remove(gen); + } + } + + /** + * Closes this manager. Idempotent. DecRefs the current snapshot and removes it if count reaches zero. + */ + @Override + public void close() { + if (closed.compareAndSet(false, true)) { + decRefAndRemove(latestCatalogSnapshot); + } + } + +} diff --git a/server/src/main/java/org/opensearch/index/engine/exec/coord/DataformatAwareCatalogSnapshot.java b/server/src/main/java/org/opensearch/index/engine/exec/coord/DataformatAwareCatalogSnapshot.java new file mode 100644 index 0000000000000..9426bbeaad47b --- /dev/null +++ b/server/src/main/java/org/opensearch/index/engine/exec/coord/DataformatAwareCatalogSnapshot.java @@ -0,0 +1,208 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +package org.opensearch.index.engine.exec.coord; + +import org.opensearch.common.annotation.ExperimentalApi; +import org.opensearch.common.io.stream.BytesStreamOutput; +import org.opensearch.core.common.bytes.BytesReference; +import org.opensearch.core.common.io.stream.BytesStreamInput; +import org.opensearch.core.common.io.stream.StreamInput; +import org.opensearch.core.common.io.stream.StreamOutput; +import org.opensearch.index.engine.dataformat.DataFormat; +import org.opensearch.index.engine.exec.Segment; +import org.opensearch.index.engine.exec.WriterFileSet; + +import java.io.IOException; +import java.util.ArrayList; +import java.util.Base64; +import java.util.Collection; +import java.util.Collections; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.Function; + +/** + * Concrete implementation of {@link CatalogSnapshot} for the composite multi-format engine. + * Holds segments grouped by data format, supports searchable file lookups across formats, + * and tracks snapshot metadata including user data and writer generation. + */ +@ExperimentalApi +public class DataformatAwareCatalogSnapshot extends CatalogSnapshot { + + private final long id; + private final List segments; + private final long lastWriterGeneration; + private Map userData; + private final AtomicBoolean closed = new AtomicBoolean(false); + + /** + * Constructs a new DataformatAwareCatalogSnapshot. + * + * @param id the unique snapshot identifier + * @param generation the monotonically increasing generation number + * @param version the schema version for serialization compatibility + * @param segments the list of segments in this snapshot + * @param lastWriterGeneration the generation of the last writer that contributed to this snapshot + * @param userData user-defined metadata key-value pairs + */ + DataformatAwareCatalogSnapshot( + long id, + long generation, + long version, + List segments, + long lastWriterGeneration, + Map userData + ) { + super("dataformat_aware_catalog_snapshot", generation, version); + this.id = id; + this.segments = Collections.unmodifiableList(new ArrayList<>(segments)); + this.lastWriterGeneration = lastWriterGeneration; + this.userData = Map.copyOf(userData); + } + + /** + * Constructs a DataformatAwareCatalogSnapshot from a {@link StreamInput}. + * + * @param in the stream input to read from + * @param directoryResolver function that maps a data format name to its directory path + * @throws IOException if an I/O error occurs + */ + DataformatAwareCatalogSnapshot(StreamInput in, Function directoryResolver) throws IOException { + super(in); + + this.userData = in.readMap(StreamInput::readString, StreamInput::readString); + + this.id = in.readLong(); + this.lastWriterGeneration = in.readLong(); + + int segmentCount = in.readVInt(); + List segmentList = new ArrayList<>(segmentCount); + for (int i = 0; i < segmentCount; i++) { + segmentList.add(new Segment(in, directoryResolver)); + } + this.segments = Collections.unmodifiableList(segmentList); + } + + @Override + public long getId() { + return id; + } + + @Override + public List getSegments() { + return segments; + } + + @Override + public Collection getSearchableFiles(String dataFormat) { + List result = new ArrayList<>(); + for (Segment segment : segments) { + WriterFileSet writerFileSet = segment.dfGroupedSearchableFiles().get(dataFormat); + if (writerFileSet != null) { + result.add(writerFileSet); + } + } + return result; + } + + @Override + public Set getDataFormats() { + Set formats = new HashSet<>(); + for (Segment segment : segments) { + formats.addAll(segment.dfGroupedSearchableFiles().keySet()); + } + return formats; + } + + @Override + public long getLastWriterGeneration() { + return lastWriterGeneration; + } + + @Override + public Map getUserData() { + return userData; + } + + @Override + public void setUserData(Map userData) { + this.userData = Map.copyOf(userData); + } + + @Override + public String serializeToString() throws IOException { + try (BytesStreamOutput out = new BytesStreamOutput()) { + this.writeTo(out); + return Base64.getEncoder().encodeToString(BytesReference.toBytes(out.bytes())); + } + } + + /** + * Deserializes a {@link DataformatAwareCatalogSnapshot} from a Base64-encoded binary string. + * + * @param serializedData the Base64 string produced by {@link #serializeToString()} + * @param directoryResolver function that maps a data format name to its directory path + * @return a reconstructed {@link DataformatAwareCatalogSnapshot} + * @throws IOException if the data is malformed or missing required fields + */ + public static DataformatAwareCatalogSnapshot deserializeFromString(String serializedData, Function directoryResolver) + throws IOException { + if (serializedData == null || serializedData.isEmpty()) { + throw new IOException("Cannot deserialize DataformatAwareCatalogSnapshot: input is null or empty"); + } + try { + byte[] bytes = Base64.getDecoder().decode(serializedData); + try (BytesStreamInput in = new BytesStreamInput(bytes)) { + return new DataformatAwareCatalogSnapshot(in, directoryResolver); + } + } catch (IOException e) { + throw e; + } catch (Exception e) { + throw new IOException("Failed to deserialize DataformatAwareCatalogSnapshot: " + e.getMessage(), e); + } + } + + @Override + public void writeTo(StreamOutput out) throws IOException { + super.writeTo(out); + out.writeMap(userData, StreamOutput::writeString, StreamOutput::writeString); + out.writeLong(id); + out.writeLong(lastWriterGeneration); + out.writeVInt(segments.size()); + for (Segment seg : segments) { + seg.writeTo(out); + } + } + + @Override + public DataformatAwareCatalogSnapshot clone() { + return new DataformatAwareCatalogSnapshot(id, generation, version, segments, lastWriterGeneration, userData); + } + + @Override + protected void closeInternal() { + closed.set(true); + } + + /** + * Returns {@code true} if {@link #closeInternal()} has been invoked (ref count reached zero). + * This method is intended for testing only. + */ + public boolean isClosed() { + return closed.get(); + } + + @Override + public Object getReader(DataFormat dataFormat) { + throw new UnsupportedOperationException("Not implemented"); + } +} diff --git a/server/src/main/java/org/opensearch/index/engine/exec/coord/SegmentInfosCatalogSnapshot.java b/server/src/main/java/org/opensearch/index/engine/exec/coord/SegmentInfosCatalogSnapshot.java new file mode 100644 index 0000000000000..82cf48f07d804 --- /dev/null +++ b/server/src/main/java/org/opensearch/index/engine/exec/coord/SegmentInfosCatalogSnapshot.java @@ -0,0 +1,145 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +package org.opensearch.index.engine.exec.coord; + +import org.apache.lucene.index.SegmentInfos; +import org.apache.lucene.store.BufferedChecksumIndexInput; +import org.apache.lucene.store.ByteBuffersDataOutput; +import org.apache.lucene.store.ByteBuffersIndexOutput; +import org.opensearch.common.annotation.ExperimentalApi; +import org.opensearch.common.lucene.store.ByteArrayIndexInput; +import org.opensearch.core.common.io.stream.StreamInput; +import org.opensearch.core.common.io.stream.StreamOutput; +import org.opensearch.index.engine.dataformat.DataFormat; +import org.opensearch.index.engine.exec.Segment; +import org.opensearch.index.engine.exec.WriterFileSet; + +import java.io.IOException; +import java.util.Collection; +import java.util.List; +import java.util.Map; +import java.util.Set; + +/** + * A thin adapter that wraps Lucene's {@link SegmentInfos} as a {@link CatalogSnapshot}. + * Used by {@code InternalEngine} (the standard single-format Lucene engine) to participate + * in the {@link CatalogSnapshot} abstraction without requiring composite engine infrastructure. + * + *

Multi-format methods ({@link #getSegments()}, {@link #getSearchableFiles(String)}, + * {@link #getDataFormats()}, {@link #serializeToString()}) throw {@link UnsupportedOperationException} + * since Lucene-only engines do not use composite segments.

+ */ +@ExperimentalApi +public class SegmentInfosCatalogSnapshot extends CatalogSnapshot { + + private static final String CATALOG_SNAPSHOT_KEY = "_segment_infos_catalog_snapshot_"; + + private final SegmentInfos segmentInfos; + + /** + * Constructs a new SegmentInfosCatalogSnapshot wrapping the given SegmentInfos. + * + * @param segmentInfos the Lucene SegmentInfos to wrap + */ + public SegmentInfosCatalogSnapshot(SegmentInfos segmentInfos) { + super(CATALOG_SNAPSHOT_KEY + segmentInfos.getGeneration(), segmentInfos.getGeneration(), segmentInfos.getVersion()); + this.segmentInfos = segmentInfos; + } + + /** + * Constructs a SegmentInfosCatalogSnapshot from a {@link StreamInput} by deserializing the + * SegmentInfos binary representation. + * + * @param in the stream input to read from + * @throws IOException if an I/O error occurs + */ + public SegmentInfosCatalogSnapshot(StreamInput in) throws IOException { + super(in); + byte[] segmentInfosBytes = in.readByteArray(); + this.segmentInfos = SegmentInfos.readCommit( + null, + new BufferedChecksumIndexInput(new ByteArrayIndexInput("SegmentInfos", segmentInfosBytes)), + 0L + ); + } + + /** + * Returns the wrapped Lucene SegmentInfos instance. + * + * @return the SegmentInfos + */ + public SegmentInfos getSegmentInfos() { + return segmentInfos; + } + + @Override + public long getId() { + return generation; + } + + @Override + public Map getUserData() { + return segmentInfos.getUserData(); + } + + @Override + public long getLastWriterGeneration() { + return -1; + } + + @Override + public List getSegments() { + throw new UnsupportedOperationException("SegmentInfosCatalogSnapshot does not support getSegments()"); + } + + @Override + public Collection getSearchableFiles(String dataFormat) { + throw new UnsupportedOperationException("SegmentInfosCatalogSnapshot does not support getSearchableFiles()"); + } + + @Override + public Set getDataFormats() { + throw new UnsupportedOperationException("SegmentInfosCatalogSnapshot does not support getDataFormats()"); + } + + @Override + public String serializeToString() throws IOException { + throw new UnsupportedOperationException("SegmentInfosCatalogSnapshot does not support serializeToString()"); + } + + @Override + public void writeTo(StreamOutput out) throws IOException { + super.writeTo(out); + ByteBuffersDataOutput buffer = new ByteBuffersDataOutput(); + try (ByteBuffersIndexOutput indexOutput = new ByteBuffersIndexOutput(buffer, "", null)) { + segmentInfos.write(indexOutput); + } + out.writeByteArray(buffer.toArrayCopy()); + } + + @Override + public void setUserData(Map userData) { + // No-op for SegmentInfosCatalogSnapshot + } + + @Override + public Object getReader(DataFormat dataFormat) { + throw new UnsupportedOperationException("SegmentInfosCatalogSnapshot does not support getReader()"); + } + + @Override + protected void closeInternal() { + // No resources to release for SegmentInfos wrapper. + } + + @Override + public SegmentInfosCatalogSnapshot clone() { + return new SegmentInfosCatalogSnapshot(segmentInfos); + } +} diff --git a/server/src/main/java/org/opensearch/index/engine/exec/coord/package-info.java b/server/src/main/java/org/opensearch/index/engine/exec/coord/package-info.java new file mode 100644 index 0000000000000..53ae20e9b9aef --- /dev/null +++ b/server/src/main/java/org/opensearch/index/engine/exec/coord/package-info.java @@ -0,0 +1,14 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +/** + * Coordination layer for the composite multi-format engine. + * Contains the CatalogSnapshotManager and related snapshot implementations + * for managing immutable point-in-time views of index data across multiple data formats. + */ +package org.opensearch.index.engine.exec.coord; diff --git a/server/src/test/java/org/opensearch/index/engine/dataformat/DataFormatPluginTests.java b/server/src/test/java/org/opensearch/index/engine/dataformat/DataFormatPluginTests.java index c6d8a01d448e8..4c5191ac1af2c 100644 --- a/server/src/test/java/org/opensearch/index/engine/dataformat/DataFormatPluginTests.java +++ b/server/src/test/java/org/opensearch/index/engine/dataformat/DataFormatPluginTests.java @@ -10,6 +10,7 @@ import org.opensearch.Version; import org.opensearch.cluster.metadata.IndexMetadata; +import org.opensearch.common.concurrent.GatedCloseable; import org.opensearch.common.settings.Settings; import org.opensearch.core.index.shard.ShardId; import org.opensearch.index.IndexSettings; @@ -23,6 +24,9 @@ import org.opensearch.index.engine.dataformat.stub.MockReaderManager; import org.opensearch.index.engine.exec.Segment; import org.opensearch.index.engine.exec.WriterFileSet; +import org.opensearch.index.engine.exec.coord.CatalogSnapshot; +import org.opensearch.index.engine.exec.coord.CatalogSnapshotManager; +import org.opensearch.index.engine.exec.coord.DataformatAwareCatalogSnapshot; import org.opensearch.index.mapper.MappedFieldType; import org.opensearch.index.mapper.MapperService; import org.opensearch.index.shard.ShardPath; @@ -30,7 +34,6 @@ import java.io.IOException; import java.nio.file.Path; -import java.util.Collection; import java.util.List; import java.util.Map; import java.util.Optional; @@ -230,14 +233,8 @@ public void testRefreshInput() { /** * Search holds snapshot alive while refresh replaces it. - *

- * Timeline: - * 1. new s1 → refcount = 1 (construction) - * 2. setLatestSnapshot(s1) → refcount = 1 (engine takes over construction ref) - * 3. acquireReader() → refcount = 2 (search adds ref) - * 4. setLatestSnapshot(s2) → s1 refcount = 1 (engine releases s1) - * 5. readerManager.onDeleted(s1) → reader closed, but s1 alive (search ref) - * 6. compositeReader.close() → s1 refcount = 0 → dead + * CatalogSnapshotManager handles ref counting: acquireReader increments, + * commitNewSnapshot replaces the latest, and closing the reader releases the old snapshot. */ public void testSearchHoldsSnapshotAliveWhileRefreshDeletesFiles() throws IOException { MockDataFormat format = new MockDataFormat(); @@ -253,20 +250,24 @@ public void testSearchHoldsSnapshotAliveWhileRefreshDeletesFiles() throws IOExce w1.close(); RefreshResult rr1 = indexEngine.refresh(RefreshInput.builder().addWriterFileSet(fs1).build()); - MockCatalogSnapshot snapshot1 = new MockCatalogSnapshot(1L, rr1.refreshedSegments(), format); + + CatalogSnapshotManager manager = new CatalogSnapshotManager(1L, 1L, 0L, rr1.refreshedSegments(), 1L, Map.of()); MockReaderManager readerManager = new MockReaderManager(format.name()); - readerManager.afterRefresh(true, snapshot1); + try (GatedCloseable ref = manager.acquireSnapshot()) { + readerManager.afterRefresh(true, ref.get()); + } - DataFormatAwareEngine dataFormatAwareEngine = new DataFormatAwareEngine(Map.of(format, readerManager)); - dataFormatAwareEngine.setLatestSnapshot(snapshot1); // takes over construction ref, refcount: 1 + DataFormatAwareEngine dataFormatAwareEngine = new DataFormatAwareEngine(Map.of(format, readerManager), manager); - // Search acquires reader — refcount: 2 + // Search acquires reader on snapshot1 — holds a ref var dataFormatAwareReader = dataFormatAwareEngine.acquireReader(); + CatalogSnapshot snapshot1 = dataFormatAwareReader.get().catalogSnapshot(); MockReader searchReader = (MockReader) dataFormatAwareReader.get().reader(format); assertEquals(1, searchReader.totalRows); + assertEquals(1L, snapshot1.getGeneration()); - // New refresh arrives — setLatestSnapshot(s2) decRefs s1 → refcount: 1 + // New refresh arrives — commit replaces snapshot Writer w2 = indexEngine.createWriter(2L); MockDocumentInput d2 = indexEngine.newDocumentInput(); d2.addField(mock(MappedFieldType.class), "Bob"); @@ -276,27 +277,41 @@ public void testSearchHoldsSnapshotAliveWhileRefreshDeletesFiles() throws IOExce w2.close(); RefreshResult rr2 = indexEngine.refresh(RefreshInput.builder().addWriterFileSet(fs1).addWriterFileSet(fs2).build()); - MockCatalogSnapshot snapshot2 = new MockCatalogSnapshot(2L, rr2.refreshedSegments(), format); - readerManager.afterRefresh(true, snapshot2); - dataFormatAwareEngine.setLatestSnapshot(snapshot2); // s1 refcount: 1 (only search ref) + manager.commitNewSnapshot(rr2.refreshedSegments()); + + try (GatedCloseable ref = manager.acquireSnapshot()) { + readerManager.afterRefresh(true, ref.get()); + } - // Old snapshot deleted from reader manager — reader closes - readerManager.onDeleted(snapshot1); - assertTrue("Reader for snapshot1 closed in reader manager", searchReader.closed); + // Snapshot1 still alive — search reader still works because the ref is held + assertFalse("Snapshot1 should still be alive while search holds ref", ((DataformatAwareCatalogSnapshot) snapshot1).isClosed()); + assertEquals(1, searchReader.totalRows); + assertSame(snapshot1, dataFormatAwareReader.get().catalogSnapshot()); - // But snapshot1 still alive — search holds the last ref - assertTrue("Snapshot1 alive while search holds ref", snapshot1.tryIncRef()); - snapshot1.decRef(); // undo probe + // New acquireSnapshot returns snapshot2, not snapshot1 + try (GatedCloseable ref = manager.acquireSnapshot()) { + assertEquals(2L, ref.get().getGeneration()); + assertNotSame(snapshot1, ref.get()); + } - // Search completes — s1 refcount: 0 → dead + // Search completes — releases the old snapshot ref dataFormatAwareReader.close(); - assertFalse("Snapshot1 dead after search releases", snapshot1.tryIncRef()); - // Snapshot 2 still works + // Snapshot1 is now dead + assertTrue( + "Snapshot1 should be closed after search releases the last ref", + ((DataformatAwareCatalogSnapshot) snapshot1).isClosed() + ); + + // Snapshot1 is now dead — tryIncRef would fail (verified via new acquire returning snapshot2) + // Snapshot 2 works try (var cr2 = dataFormatAwareEngine.acquireReader()) { MockReader r2 = (MockReader) cr2.get().reader(format); assertEquals(2, r2.totalRows); + assertEquals(2L, cr2.get().catalogSnapshot().getGeneration()); } + + manager.close(); } /** @@ -328,24 +343,15 @@ public Set supportedFields() { WriterFileSet wfs1 = WriterFileSet.builder().directory(dir).writerGeneration(1L).addFile("data.parquet").addNumRows(10).build(); WriterFileSet wfs2 = WriterFileSet.builder().directory(dir).writerGeneration(1L).addFile("data.lucene").addNumRows(10).build(); Segment seg = Segment.builder(0L).addSearchableFiles(format1, wfs1).addSearchableFiles(format2, wfs2).build(); - MockCatalogSnapshot snapshot = new MockCatalogSnapshot(1L, List.of(seg), format1) { - @Override - public Collection getSearchableFiles(String dataFormat) { - if ("mock-lucene".equals(dataFormat)) return List.of(wfs2); - return super.getSearchableFiles(dataFormat); - } - @Override - public Set getDataFormats() { - return Set.of(format1.name(), format2.name()); - } - }; + CatalogSnapshotManager manager = new CatalogSnapshotManager(1L, 1L, 0L, List.of(seg), 1L, Map.of()); - rm1.afterRefresh(true, snapshot); - rm2.afterRefresh(true, snapshot); + try (GatedCloseable ref = manager.acquireSnapshot()) { + rm1.afterRefresh(true, ref.get()); + rm2.afterRefresh(true, ref.get()); + } - DataFormatAwareEngine dataFormatAwareEngine = new DataFormatAwareEngine(Map.of(format1, rm1, format2, rm2)); - dataFormatAwareEngine.setLatestSnapshot(snapshot); + DataFormatAwareEngine dataFormatAwareEngine = new DataFormatAwareEngine(Map.of(format1, rm1, format2, rm2), manager); try (var cr = dataFormatAwareEngine.acquireReader()) { MockReader r1 = (MockReader) cr.get().reader(format1); @@ -357,6 +363,8 @@ public Set getDataFormats() { assertTrue(r1.fileNames.contains("data.parquet")); assertTrue(r2.fileNames.contains("data.lucene")); } + + manager.close(); } /** diff --git a/server/src/test/java/org/opensearch/index/engine/dataformat/merge/MergeTests.java b/server/src/test/java/org/opensearch/index/engine/dataformat/merge/MergeTests.java index fb7cf71caa84f..9444d0d6d11f8 100644 --- a/server/src/test/java/org/opensearch/index/engine/dataformat/merge/MergeTests.java +++ b/server/src/test/java/org/opensearch/index/engine/dataformat/merge/MergeTests.java @@ -16,10 +16,10 @@ import org.opensearch.index.engine.dataformat.DataFormat; import org.opensearch.index.engine.dataformat.MergeResult; import org.opensearch.index.engine.dataformat.stub.MockDataFormat; -import org.opensearch.index.engine.exec.CatalogSnapshot; import org.opensearch.index.engine.exec.Indexer; import org.opensearch.index.engine.exec.Segment; import org.opensearch.index.engine.exec.WriterFileSet; +import org.opensearch.index.engine.exec.coord.CatalogSnapshot; import org.opensearch.test.OpenSearchTestCase; import java.nio.file.Path; diff --git a/server/src/test/java/org/opensearch/index/engine/dataformat/stub/MockCatalogSnapshot.java b/server/src/test/java/org/opensearch/index/engine/dataformat/stub/MockCatalogSnapshot.java index d5af44b775abf..9d619d95ccbcb 100644 --- a/server/src/test/java/org/opensearch/index/engine/dataformat/stub/MockCatalogSnapshot.java +++ b/server/src/test/java/org/opensearch/index/engine/dataformat/stub/MockCatalogSnapshot.java @@ -8,11 +8,13 @@ package org.opensearch.index.engine.dataformat.stub; +import org.opensearch.core.common.io.stream.StreamOutput; import org.opensearch.index.engine.dataformat.DataFormat; -import org.opensearch.index.engine.exec.CatalogSnapshot; import org.opensearch.index.engine.exec.Segment; import org.opensearch.index.engine.exec.WriterFileSet; +import org.opensearch.index.engine.exec.coord.CatalogSnapshot; +import java.io.IOException; import java.util.ArrayList; import java.util.Collection; import java.util.List; @@ -73,16 +75,23 @@ public String serializeToString() { } @Override - public void setCatalogSnapshotMap(Map map) {} - - @Override - public void setUserData(Map userData, boolean b) {} + public void setUserData(Map userData) {} @Override public Object getReader(DataFormat dataFormat) { return null; } + @Override + public CatalogSnapshot clone() { + return new MockCatalogSnapshot(generation, segments, format); + } + + @Override + public void writeTo(StreamOutput out) throws IOException { + super.writeTo(out); + } + @Override protected void closeInternal() {} } diff --git a/server/src/test/java/org/opensearch/index/engine/dataformat/stub/MockReaderManager.java b/server/src/test/java/org/opensearch/index/engine/dataformat/stub/MockReaderManager.java index bfb16bbf2d329..7f628e73d26fa 100644 --- a/server/src/test/java/org/opensearch/index/engine/dataformat/stub/MockReaderManager.java +++ b/server/src/test/java/org/opensearch/index/engine/dataformat/stub/MockReaderManager.java @@ -8,9 +8,9 @@ package org.opensearch.index.engine.dataformat.stub; -import org.opensearch.index.engine.exec.CatalogSnapshot; import org.opensearch.index.engine.exec.EngineReaderManager; import org.opensearch.index.engine.exec.WriterFileSet; +import org.opensearch.index.engine.exec.coord.CatalogSnapshot; import java.util.ArrayList; import java.util.Collection; diff --git a/server/src/test/java/org/opensearch/index/engine/exec/SegmentTests.java b/server/src/test/java/org/opensearch/index/engine/exec/SegmentTests.java new file mode 100644 index 0000000000000..d5afc4257c4a1 --- /dev/null +++ b/server/src/test/java/org/opensearch/index/engine/exec/SegmentTests.java @@ -0,0 +1,82 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +package org.opensearch.index.engine.exec; + +import org.opensearch.core.common.io.stream.NamedWriteableRegistry; +import org.opensearch.test.OpenSearchTestCase; + +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.Map; +import java.util.Set; + +/** + * Tests for {@link Segment} serialization. + */ +public class SegmentTests extends OpenSearchTestCase { + + private static final String TEST_DIRECTORY = "/tmp/test-segment"; + + public void testCopyWriteable() throws Exception { + Segment original = randomSegment(); + Segment copy = copyWriteable( + original, + new NamedWriteableRegistry(Collections.emptyList()), + in -> new Segment(in, key -> TEST_DIRECTORY) + ); + assertEquals(original, copy); + } + + public void testCopyWriteableEmpty() throws Exception { + Segment empty = new Segment(0L, Map.of()); + Segment copy = copyWriteable( + empty, + new NamedWriteableRegistry(Collections.emptyList()), + in -> new Segment(in, key -> TEST_DIRECTORY) + ); + assertEquals(empty, copy); + } + + public void testCopyWriteableMultiFormat() throws Exception { + Map dfGrouped = new HashMap<>(); + dfGrouped.put("lucene", randomWriterFileSet("lucene")); + dfGrouped.put("parquet", randomWriterFileSet("parquet")); + Segment original = new Segment(randomNonNegativeLong(), dfGrouped); + + Segment copy = copyWriteable( + original, + new NamedWriteableRegistry(Collections.emptyList()), + in -> new Segment(in, key -> TEST_DIRECTORY) + ); + assertEquals(original, copy); + assertEquals(2, copy.dfGroupedSearchableFiles().size()); + } + + // --- helpers --- + + private WriterFileSet randomWriterFileSet(String format) { + int fileCount = randomIntBetween(1, 5); + Set files = new HashSet<>(); + String[] extensions = "lucene".equals(format) ? new String[] { "cfs", "si", "dat" } : new String[] { "parquet" }; + for (int i = 0; i < fileCount; i++) { + files.add(randomAlphaOfLength(6) + "." + randomFrom(extensions)); + } + return new WriterFileSet(TEST_DIRECTORY, randomNonNegativeLong(), files, randomIntBetween(0, 10000)); + } + + private Segment randomSegment() { + Map dfGrouped = new HashMap<>(); + for (int i = 0; i < randomIntBetween(1, 3); i++) { + String format = randomFrom("lucene", "parquet"); + dfGrouped.put(format, randomWriterFileSet(format)); + } + return new Segment(randomNonNegativeLong(), dfGrouped); + } +} diff --git a/server/src/test/java/org/opensearch/index/engine/exec/WriterFileSetTests.java b/server/src/test/java/org/opensearch/index/engine/exec/WriterFileSetTests.java new file mode 100644 index 0000000000000..2eb0d82b92728 --- /dev/null +++ b/server/src/test/java/org/opensearch/index/engine/exec/WriterFileSetTests.java @@ -0,0 +1,63 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +package org.opensearch.index.engine.exec; + +import org.opensearch.core.common.io.stream.NamedWriteableRegistry; +import org.opensearch.test.OpenSearchTestCase; + +import java.util.Collections; +import java.util.HashSet; +import java.util.Set; + +/** + * Tests for {@link WriterFileSet}. + */ +public class WriterFileSetTests extends OpenSearchTestCase { + + public void testCopyWriteable() throws Exception { + WriterFileSet original = randomWriterFileSet(); + String directory = original.directory(); + WriterFileSet copy = copyWriteable( + original, + new NamedWriteableRegistry(Collections.emptyList()), + in -> new WriterFileSet(in, directory) + ); + assertEquals(original, copy); + } + + public void testDirectoryNotSerialized() throws Exception { + String originalDirectory = "/tmp/original"; + String differentDirectory = "/tmp/different"; + WriterFileSet original = new WriterFileSet(originalDirectory, 1L, Set.of("a.dat"), 10); + + WriterFileSet deserialized = copyWriteable( + original, + new NamedWriteableRegistry(Collections.emptyList()), + in -> new WriterFileSet(in, differentDirectory) + ); + + assertEquals(differentDirectory, deserialized.directory()); + assertNotEquals(originalDirectory, deserialized.directory()); + assertEquals(original.writerGeneration(), deserialized.writerGeneration()); + assertEquals(original.files(), deserialized.files()); + assertEquals(original.numRows(), deserialized.numRows()); + } + + // --- helpers --- + + private WriterFileSet randomWriterFileSet() { + String directory = "/tmp/" + randomAlphaOfLength(8); + int fileCount = randomIntBetween(1, 5); + Set files = new HashSet<>(); + for (int i = 0; i < fileCount; i++) { + files.add(randomAlphaOfLength(6) + "." + randomFrom("cfs", "si", "dat", "parquet")); + } + return new WriterFileSet(directory, randomNonNegativeLong(), files, randomIntBetween(0, 10000)); + } +} diff --git a/server/src/test/java/org/opensearch/index/engine/exec/coord/CatalogSnapshotManagerTests.java b/server/src/test/java/org/opensearch/index/engine/exec/coord/CatalogSnapshotManagerTests.java new file mode 100644 index 0000000000000..5814def67f666 --- /dev/null +++ b/server/src/test/java/org/opensearch/index/engine/exec/coord/CatalogSnapshotManagerTests.java @@ -0,0 +1,297 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +package org.opensearch.index.engine.exec.coord; + +import org.opensearch.common.concurrent.GatedCloseable; +import org.opensearch.index.engine.exec.Segment; +import org.opensearch.index.engine.exec.WriterFileSet; +import org.opensearch.test.OpenSearchTestCase; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; + +/** + * Tests for {@link CatalogSnapshotManager}. + */ +public class CatalogSnapshotManagerTests extends OpenSearchTestCase { + + public void testCommitProducesCorrectNewSnapshot() throws Exception { + for (int iter = 0; iter < 100; iter++) { + CatalogSnapshotManager manager = createRandomManager(); + try { + long previousGeneration; + Set seenIds = new HashSet<>(); + try (GatedCloseable ref = manager.acquireSnapshot()) { + previousGeneration = ref.get().getGeneration(); + seenIds.add(ref.get().getId()); + } + + int numCommits = randomIntBetween(1, 10); + for (int c = 0; c < numCommits; c++) { + List newSegments = randomSegments(); + manager.commitNewSnapshot(newSegments); + + try (GatedCloseable ref = manager.acquireSnapshot()) { + assertEquals(previousGeneration + 1, ref.get().getGeneration()); + assertTrue(seenIds.add(ref.get().getId())); + assertEquals(newSegments, ref.get().getSegments()); + previousGeneration = ref.get().getGeneration(); + } + } + } finally { + manager.close(); + } + } + } + + public void testUserDataPreservationOnCommit() throws Exception { + for (int iter = 0; iter < 100; iter++) { + Map initialUserData = randomUserData(randomIntBetween(1, 5)); + long initGen = randomIntBetween(0, 100); + CatalogSnapshotManager manager = new CatalogSnapshotManager( + randomNonNegativeLong(), + initGen, + randomNonNegativeLong(), + randomSegments(), + randomNonNegativeLong(), + initialUserData + ); + try { + manager.commitNewSnapshot(randomSegments()); + try (GatedCloseable ref = manager.acquireSnapshot()) { + assertEquals(initialUserData, ref.get().getUserData()); + } + + manager.commitNewSnapshot(randomSegments()); + try (GatedCloseable ref = manager.acquireSnapshot()) { + assertEquals(initialUserData, ref.get().getUserData()); + } + } finally { + manager.close(); + } + } + } + + public void testReferenceCountingLifecycle() throws Exception { + for (int iter = 0; iter < 100; iter++) { + long initGen = randomIntBetween(0, 100); + CatalogSnapshotManager manager = new CatalogSnapshotManager( + randomNonNegativeLong(), + initGen, + randomNonNegativeLong(), + randomSegments(), + randomNonNegativeLong(), + Collections.emptyMap() + ); + + CatalogSnapshot initialSnapshot; + try (GatedCloseable ref = manager.acquireSnapshot()) { + initialSnapshot = ref.get(); + assertEquals(2, initialSnapshot.refCount()); + } + assertEquals(1, initialSnapshot.refCount()); + + manager.commitNewSnapshot(randomSegments()); + assertEquals(0, initialSnapshot.refCount()); + + int numCommits = randomIntBetween(1, 8); + for (int c = 0; c < numCommits; c++) { + CatalogSnapshot prev; + try (GatedCloseable ref = manager.acquireSnapshot()) { + prev = ref.get(); + assertEquals(2, prev.refCount()); + } + assertEquals(1, prev.refCount()); + manager.commitNewSnapshot(randomSegments()); + assertEquals(0, prev.refCount()); + } + + CatalogSnapshot finalSnapshot; + try (GatedCloseable ref = manager.acquireSnapshot()) { + finalSnapshot = ref.get(); + assertEquals(2, finalSnapshot.refCount()); + } + assertEquals(1, finalSnapshot.refCount()); + manager.close(); + assertEquals(0, finalSnapshot.refCount()); + } + } + + public void testAcquireAndReleaseViaGatedCloseable() throws Exception { + for (int iter = 0; iter < 100; iter++) { + CatalogSnapshotManager manager = createRandomManager(); + try { + CatalogSnapshot currentSnap; + try (GatedCloseable initialRef = manager.acquireSnapshot()) { + currentSnap = initialRef.get(); + assertEquals(2, currentSnap.refCount()); + } + assertEquals(1, currentSnap.refCount()); + + int numAcquires = randomIntBetween(1, 5); + List> refs = new ArrayList<>(); + for (int a = 0; a < numAcquires; a++) { + refs.add(manager.acquireSnapshot()); + assertEquals(1 + (a + 1), currentSnap.refCount()); + } + for (int r = 0; r < numAcquires; r++) { + refs.get(r).close(); + assertEquals(1 + numAcquires - r - 1, currentSnap.refCount()); + } + assertEquals(1, currentSnap.refCount()); + + GatedCloseable heldRef = manager.acquireSnapshot(); + CatalogSnapshot heldSnapshot = heldRef.get(); + assertEquals(2, heldSnapshot.refCount()); + + manager.commitNewSnapshot(randomSegments()); + assertEquals(1, heldSnapshot.refCount()); + + heldRef.close(); + assertEquals(0, heldSnapshot.refCount()); + } finally { + manager.close(); + } + } + } + + public void testClosedManagerRejectsAcquisition() throws Exception { + for (int iter = 0; iter < 100; iter++) { + CatalogSnapshotManager manager = createRandomManager(); + for (int c = 0; c < randomIntBetween(0, 5); c++) { + manager.commitNewSnapshot(randomSegments()); + } + manager.close(); + expectThrows(IllegalStateException.class, manager::acquireSnapshot); + } + } + + public void testInitialSnapshotRecovery() throws Exception { + for (int iter = 0; iter < 100; iter++) { + long id = randomNonNegativeLong(); + long generation = randomIntBetween(0, 100); + long version = randomNonNegativeLong(); + long lastWriterGeneration = randomNonNegativeLong(); + List segments = randomIntBetween(1, 5) == 1 ? Collections.emptyList() : randomSegments(); + Map userData = randomUserData(randomIntBetween(0, 4)); + + CatalogSnapshotManager manager = new CatalogSnapshotManager(id, generation, version, segments, lastWriterGeneration, userData); + try (GatedCloseable ref = manager.acquireSnapshot()) { + CatalogSnapshot acquired = ref.get(); + assertEquals(id, acquired.getId()); + assertEquals(generation, acquired.getGeneration()); + assertEquals(segments, acquired.getSegments()); + assertEquals(userData, acquired.getUserData()); + assertEquals(lastWriterGeneration, acquired.getLastWriterGeneration()); + } finally { + manager.close(); + } + } + } + + public void testCloseInternalInvokedOnCommit() throws Exception { + CatalogSnapshotManager manager = createRandomManager(); + + CatalogSnapshot initialSnapshot; + try (GatedCloseable ref = manager.acquireSnapshot()) { + initialSnapshot = ref.get(); + } + assertFalse(((DataformatAwareCatalogSnapshot) initialSnapshot).isClosed()); + + manager.commitNewSnapshot(randomSegments()); + assertTrue( + "snapshot should be closed when commit replaces the last ref", + ((DataformatAwareCatalogSnapshot) initialSnapshot).isClosed() + ); + manager.close(); + } + + public void testCloseInternalInvokedOnManagerClose() throws Exception { + CatalogSnapshotManager manager = createRandomManager(); + + CatalogSnapshot snapshot; + try (GatedCloseable ref = manager.acquireSnapshot()) { + snapshot = ref.get(); + } + assertFalse(((DataformatAwareCatalogSnapshot) snapshot).isClosed()); + + manager.close(); + assertTrue("snapshot should be closed when manager releases the last ref", ((DataformatAwareCatalogSnapshot) snapshot).isClosed()); + } + + public void testCloseInternalNotInvokedWhileRefsHeld() throws Exception { + CatalogSnapshotManager manager = createRandomManager(); + + GatedCloseable heldRef = manager.acquireSnapshot(); + CatalogSnapshot heldSnapshot = heldRef.get(); + assertFalse(((DataformatAwareCatalogSnapshot) heldSnapshot).isClosed()); + + manager.commitNewSnapshot(randomSegments()); + assertFalse("snapshot should not be closed while a ref is still held", ((DataformatAwareCatalogSnapshot) heldSnapshot).isClosed()); + + heldRef.close(); + assertTrue("snapshot should be closed after the last ref is released", ((DataformatAwareCatalogSnapshot) heldSnapshot).isClosed()); + + manager.close(); + } + + // --- helpers --- + + private WriterFileSet randomWriterFileSet(String format) { + String directory = "/tmp/" + randomAlphaOfLength(8); + int fileCount = randomIntBetween(1, 5); + Set files = new HashSet<>(); + String[] extensions = "lucene".equals(format) ? new String[] { "cfs", "si", "dat" } : new String[] { "parquet" }; + for (int i = 0; i < fileCount; i++) { + files.add(randomAlphaOfLength(6) + "." + randomFrom(extensions)); + } + return new WriterFileSet(directory, randomNonNegativeLong(), files, randomIntBetween(0, 10000)); + } + + private Segment randomSegment() { + Map dfGrouped = new HashMap<>(); + for (int i = 0; i < randomIntBetween(1, 2); i++) { + String format = randomFrom("lucene", "parquet"); + dfGrouped.put(format, randomWriterFileSet(format)); + } + return new Segment(randomNonNegativeLong(), dfGrouped); + } + + private List randomSegments() { + List segments = new ArrayList<>(); + for (int i = 0; i < randomIntBetween(0, 5); i++) { + segments.add(randomSegment()); + } + return segments; + } + + private Map randomUserData(int entries) { + Map userData = new HashMap<>(); + for (int i = 0; i < entries; i++) { + userData.put(randomAlphaOfLength(5), randomAlphaOfLength(10)); + } + return userData; + } + + private CatalogSnapshotManager createRandomManager() { + return new CatalogSnapshotManager( + randomNonNegativeLong(), + randomIntBetween(0, 100), + randomNonNegativeLong(), + randomSegments(), + randomNonNegativeLong(), + Map.of() + ); + } +} diff --git a/server/src/test/java/org/opensearch/index/engine/exec/coord/DataformatAwareCatalogSnapshotTests.java b/server/src/test/java/org/opensearch/index/engine/exec/coord/DataformatAwareCatalogSnapshotTests.java new file mode 100644 index 0000000000000..3b6089ab4729e --- /dev/null +++ b/server/src/test/java/org/opensearch/index/engine/exec/coord/DataformatAwareCatalogSnapshotTests.java @@ -0,0 +1,512 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +package org.opensearch.index.engine.exec.coord; + +import org.opensearch.core.common.io.stream.NamedWriteableRegistry; +import org.opensearch.index.engine.exec.Segment; +import org.opensearch.index.engine.exec.WriterFileSet; +import org.opensearch.test.OpenSearchTestCase; + +import java.io.IOException; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; + +/** + * Tests for {@link DataformatAwareCatalogSnapshot}. + */ +public class DataformatAwareCatalogSnapshotTests extends OpenSearchTestCase { + + public void testSnapshotFieldAccessConsistency() { + for (int iter = 0; iter < 100; iter++) { + long id = randomLong(); + long generation = randomNonNegativeLong(); + long version = randomNonNegativeLong(); + List segments = randomSegments(); + long lastWriterGeneration = randomNonNegativeLong(); + Map userData = randomUserData(); + + DataformatAwareCatalogSnapshot snapshot = new DataformatAwareCatalogSnapshot( + id, + generation, + version, + segments, + lastWriterGeneration, + userData + ); + + assertEquals(id, snapshot.getId()); + assertEquals(generation, snapshot.getGeneration()); + assertEquals(version, snapshot.getVersion()); + assertEquals(segments, snapshot.getSegments()); + assertEquals(lastWriterGeneration, snapshot.getLastWriterGeneration()); + assertEquals(userData, snapshot.getUserData()); + + Set expectedFormats = new HashSet<>(); + for (Segment seg : segments) { + expectedFormats.addAll(seg.dfGroupedSearchableFiles().keySet()); + } + for (String format : expectedFormats) { + List expected = new ArrayList<>(); + for (Segment seg : segments) { + WriterFileSet wfs = seg.dfGroupedSearchableFiles().get(format); + if (wfs != null) expected.add(wfs); + } + assertEquals(expected, new ArrayList<>(snapshot.getSearchableFiles(format))); + } + + assertTrue(snapshot.getSearchableFiles("nonexistent_" + randomAlphaOfLength(5)).isEmpty()); + assertEquals(expectedFormats, snapshot.getDataFormats()); + expectThrows(UnsupportedOperationException.class, () -> snapshot.getSegments().add(randomSegment())); + } + } + + public void testSerializationRoundTrip() throws Exception { + for (int iter = 0; iter < 100; iter++) { + DataformatAwareCatalogSnapshot original = randomSnapshot(); + String serialized = original.serializeToString(); + + // Directory is not serialized; pass a placeholder for deserialization + String directory = "/tmp/deserialized"; + DataformatAwareCatalogSnapshot deserialized = DataformatAwareCatalogSnapshot.deserializeFromString( + serialized, + key -> directory + ); + assertSnapshotMetadataEqual("round-trip", original, deserialized); + + String reserialized = deserialized.serializeToString(); + DataformatAwareCatalogSnapshot deserialized2 = DataformatAwareCatalogSnapshot.deserializeFromString( + reserialized, + key -> directory + ); + assertSnapshotFieldsEqual("double round-trip", deserialized, deserialized2); + } + } + + public void testCopyWriteable() throws Exception { + String directory = "/tmp/" + randomAlphaOfLength(8); + DataformatAwareCatalogSnapshot original = randomSnapshotWithDirectory(directory); + DataformatAwareCatalogSnapshot copy = copyWriteable( + original, + new NamedWriteableRegistry(Collections.emptyList()), + in -> new DataformatAwareCatalogSnapshot(in, key -> directory) + ); + assertSnapshotFieldsEqual("copyWriteable", original, copy); + } + + public void testDeserializationRejectsInvalidInput() { + for (int iter = 0; iter < 100; iter++) { + String input = generateInvalidInput(iter); + expectThrows(IOException.class, () -> DataformatAwareCatalogSnapshot.deserializeFromString(input, key -> "/tmp/test")); + } + } + + public void testInitialRefCountIsOne() { + DataformatAwareCatalogSnapshot snapshot = randomSnapshot(); + assertEquals(1, snapshot.refCount()); + } + + public void testAcquireRefIncrementsCount() { + DataformatAwareCatalogSnapshot snapshot = randomSnapshot(); + assertEquals(1, snapshot.refCount()); + + snapshot.tryIncRef(); + assertEquals(2, snapshot.refCount()); + + snapshot.tryIncRef(); + assertEquals(3, snapshot.refCount()); + + snapshot.decRef(); + snapshot.decRef(); + snapshot.decRef(); + } + + public void testReleaseRefDecrementsAndTriggersCloseAtZero() { + DataformatAwareCatalogSnapshot snapshot = randomSnapshot(); + assertEquals(1, snapshot.refCount()); + + snapshot.tryIncRef(); + assertEquals(2, snapshot.refCount()); + + assertFalse(snapshot.decRef()); + assertEquals(1, snapshot.refCount()); + + assertTrue(snapshot.decRef()); + assertEquals(0, snapshot.refCount()); + } + + public void testTryAcquireRefSucceedsWhenOpen() { + DataformatAwareCatalogSnapshot snapshot = randomSnapshot(); + assertTrue(snapshot.tryIncRef()); + assertEquals(2, snapshot.refCount()); + + snapshot.decRef(); + snapshot.decRef(); + } + + public void testTryAcquireRefFailsWhenClosed() { + DataformatAwareCatalogSnapshot snapshot = randomSnapshot(); + assertTrue(snapshot.decRef()); + assertEquals(0, snapshot.refCount()); + + assertFalse(snapshot.tryIncRef()); + } + + public void testCloseInternalCalledOnceAtZeroRefCount() { + final java.util.concurrent.atomic.AtomicInteger closeCount = new java.util.concurrent.atomic.AtomicInteger(0); + DataformatAwareCatalogSnapshot snapshot = new DataformatAwareCatalogSnapshot(1L, 1L, 1L, List.of(), 0L, Map.of()) { + @Override + protected void closeInternal() { + closeCount.incrementAndGet(); + } + }; + + snapshot.tryIncRef(); + snapshot.tryIncRef(); + assertEquals(0, closeCount.get()); + + snapshot.decRef(); + assertEquals(0, closeCount.get()); + + snapshot.decRef(); + assertEquals(0, closeCount.get()); + + snapshot.decRef(); + assertEquals(1, closeCount.get()); + } + + public void testClonedSnapshotHasFreshRefCount() { + DataformatAwareCatalogSnapshot original = randomSnapshot(); + original.tryIncRef(); + assertEquals(2, original.refCount()); + + DataformatAwareCatalogSnapshot cloned = original.clone(); + assertEquals(1, cloned.refCount()); + assertEquals(2, original.refCount()); + + original.decRef(); + original.decRef(); + cloned.decRef(); + } + + public void testRefCounterDelegatesCloseInternalToSubclass() { + // Verifies the anonymous AbstractRefCounted bridge calls the subclass closeInternal, not a default + final List events = new ArrayList<>(); + DataformatAwareCatalogSnapshot snapshot = new DataformatAwareCatalogSnapshot(1L, 1L, 1L, List.of(), 0L, Map.of()) { + @Override + protected void closeInternal() { + events.add("subclass-closed"); + } + }; + + assertTrue(events.isEmpty()); + snapshot.decRef(); + assertEquals(List.of("subclass-closed"), events); + } + + public void testEachSnapshotHasIndependentRefCounter() { + DataformatAwareCatalogSnapshot snap1 = randomSnapshot(); + DataformatAwareCatalogSnapshot snap2 = randomSnapshot(); + + snap1.tryIncRef(); + assertEquals(2, snap1.refCount()); + assertEquals(1, snap2.refCount()); + + snap2.decRef(); + assertEquals(0, snap2.refCount()); + assertEquals(2, snap1.refCount()); + + snap1.decRef(); + snap1.decRef(); + } + + public void testRefCounterSurvivesMultipleIncDecCycles() { + DataformatAwareCatalogSnapshot snapshot = randomSnapshot(); + + for (int cycle = 0; cycle < 10; cycle++) { + int refs = randomIntBetween(1, 20); + for (int i = 0; i < refs; i++) { + snapshot.tryIncRef(); + } + assertEquals(1 + refs, snapshot.refCount()); + for (int i = 0; i < refs; i++) { + assertFalse(snapshot.decRef()); + } + assertEquals(1, snapshot.refCount()); + } + + assertTrue(snapshot.decRef()); + assertEquals(0, snapshot.refCount()); + assertFalse(snapshot.tryIncRef()); + } + + public void testCloseInternalNotCalledOnIntermediateDecRef() { + final java.util.concurrent.atomic.AtomicBoolean closed = new java.util.concurrent.atomic.AtomicBoolean(false); + DataformatAwareCatalogSnapshot snapshot = new DataformatAwareCatalogSnapshot(1L, 1L, 1L, List.of(), 0L, Map.of()) { + @Override + protected void closeInternal() { + closed.set(true); + } + }; + + // Acquire several refs + int extraRefs = randomIntBetween(2, 10); + for (int i = 0; i < extraRefs; i++) { + snapshot.tryIncRef(); + } + + // Release all but one — closeInternal must NOT fire + for (int i = 0; i < extraRefs; i++) { + snapshot.decRef(); + assertFalse("closeInternal should not fire while refs remain", closed.get()); + } + + // Release the last ref — now it fires + snapshot.decRef(); + assertTrue("closeInternal should fire when last ref is released", closed.get()); + } + + public void testDeserializedSnapshotHasIndependentRefCounter() throws Exception { + String directory = "/tmp/" + randomAlphaOfLength(8); + DataformatAwareCatalogSnapshot original = randomSnapshotWithDirectory(directory); + String serialized = original.serializeToString(); + + DataformatAwareCatalogSnapshot deserialized = DataformatAwareCatalogSnapshot.deserializeFromString(serialized, key -> directory); + + // Each has its own ref counter starting at 1 + assertEquals(1, original.refCount()); + assertEquals(1, deserialized.refCount()); + + original.tryIncRef(); + assertEquals(2, original.refCount()); + assertEquals(1, deserialized.refCount()); + + original.decRef(); + original.decRef(); + deserialized.decRef(); + } + + public void testIsClosedInitiallyFalse() { + DataformatAwareCatalogSnapshot snapshot = randomSnapshot(); + assertFalse(snapshot.isClosed()); + snapshot.decRef(); + } + + public void testIsClosedTrueAfterLastDecRef() { + DataformatAwareCatalogSnapshot snapshot = randomSnapshot(); + assertFalse(snapshot.isClosed()); + + snapshot.decRef(); + assertTrue(snapshot.isClosed()); + } + + public void testIsClosedFalseWhileRefsRemain() { + DataformatAwareCatalogSnapshot snapshot = randomSnapshot(); + int extraRefs = randomIntBetween(1, 10); + for (int i = 0; i < extraRefs; i++) { + snapshot.tryIncRef(); + } + + for (int i = 0; i < extraRefs; i++) { + snapshot.decRef(); + assertFalse("isClosed should be false while refs remain", snapshot.isClosed()); + } + + snapshot.decRef(); + assertTrue(snapshot.isClosed()); + } + + public void testIsClosedAfterCloneIndependent() { + DataformatAwareCatalogSnapshot original = randomSnapshot(); + DataformatAwareCatalogSnapshot cloned = original.clone(); + + original.decRef(); + assertTrue(original.isClosed()); + assertFalse(cloned.isClosed()); + + cloned.decRef(); + assertTrue(cloned.isClosed()); + } + + public void testTryIncRefFailsAfterClosed() { + DataformatAwareCatalogSnapshot snapshot = randomSnapshot(); + snapshot.decRef(); + assertTrue(snapshot.isClosed()); + assertFalse(snapshot.tryIncRef()); + } + + // --- helpers --- + + private WriterFileSet randomWriterFileSet(String format) { + String directory = "/tmp/" + randomAlphaOfLength(8); + long writerGeneration = randomNonNegativeLong(); + int fileCount = randomIntBetween(1, 5); + Set files = new HashSet<>(); + String[] extensions = "lucene".equals(format) ? new String[] { "cfs", "si", "dat" } : new String[] { "parquet" }; + for (int i = 0; i < fileCount; i++) { + files.add(randomAlphaOfLength(6) + "." + randomFrom(extensions)); + } + return new WriterFileSet(directory, writerGeneration, files, randomIntBetween(0, 10000)); + } + + private Segment randomSegment() { + long generation = randomNonNegativeLong(); + int formatCount = randomIntBetween(1, 2); + Map dfGrouped = new HashMap<>(); + for (int i = 0; i < formatCount; i++) { + String format = randomFrom("lucene", "parquet"); + dfGrouped.put(format, randomWriterFileSet(format)); + } + return new Segment(generation, dfGrouped); + } + + private List randomSegments() { + int count = randomIntBetween(0, 5); + List segments = new ArrayList<>(); + for (int i = 0; i < count; i++) { + segments.add(randomSegment()); + } + return segments; + } + + private Map randomUserData() { + int entries = randomIntBetween(0, 4); + Map userData = new HashMap<>(); + for (int i = 0; i < entries; i++) { + userData.put(randomAlphaOfLength(5), randomAlphaOfLength(10)); + } + return userData; + } + + private DataformatAwareCatalogSnapshot randomSnapshot() { + return new DataformatAwareCatalogSnapshot( + randomLong(), + randomNonNegativeLong(), + randomNonNegativeLong(), + randomSegments(), + randomNonNegativeLong(), + randomUserData() + ); + } + + private DataformatAwareCatalogSnapshot randomSnapshotWithDirectory(String directory) { + return new DataformatAwareCatalogSnapshot( + randomLong(), + randomNonNegativeLong(), + randomNonNegativeLong(), + randomSegmentsWithDirectory(directory), + randomNonNegativeLong(), + randomUserData() + ); + } + + private List randomSegmentsWithDirectory(String directory) { + int count = randomIntBetween(0, 5); + List segments = new ArrayList<>(); + for (int i = 0; i < count; i++) { + segments.add(randomSegmentWithDirectory(directory)); + } + return segments; + } + + private Segment randomSegmentWithDirectory(String directory) { + long generation = randomNonNegativeLong(); + int formatCount = randomIntBetween(1, 2); + Map dfGrouped = new HashMap<>(); + for (int i = 0; i < formatCount; i++) { + String format = randomFrom("lucene", "parquet"); + dfGrouped.put(format, randomWriterFileSetWithDirectory(format, directory)); + } + return new Segment(generation, dfGrouped); + } + + private WriterFileSet randomWriterFileSetWithDirectory(String format, String directory) { + long writerGeneration = randomNonNegativeLong(); + int fileCount = randomIntBetween(1, 5); + Set files = new HashSet<>(); + String[] extensions = "lucene".equals(format) ? new String[] { "cfs", "si", "dat" } : new String[] { "parquet" }; + for (int i = 0; i < fileCount; i++) { + files.add(randomAlphaOfLength(6) + "." + randomFrom(extensions)); + } + return new WriterFileSet(directory, writerGeneration, files, randomIntBetween(0, 10000)); + } + + private void assertSnapshotFieldsEqual(String context, DataformatAwareCatalogSnapshot expected, DataformatAwareCatalogSnapshot actual) { + assertEquals(context + ": id", expected.getId(), actual.getId()); + assertEquals(context + ": generation", expected.getGeneration(), actual.getGeneration()); + assertEquals(context + ": version", expected.getVersion(), actual.getVersion()); + assertEquals(context + ": segments", expected.getSegments(), actual.getSegments()); + assertEquals(context + ": lastWriterGeneration", expected.getLastWriterGeneration(), actual.getLastWriterGeneration()); + assertEquals(context + ": userData", expected.getUserData(), actual.getUserData()); + } + + /** + * Asserts metadata equality between two snapshots, ignoring directory (which is not serialized). + */ + private void assertSnapshotMetadataEqual( + String context, + DataformatAwareCatalogSnapshot expected, + DataformatAwareCatalogSnapshot actual + ) { + assertEquals(context + ": id", expected.getId(), actual.getId()); + assertEquals(context + ": generation", expected.getGeneration(), actual.getGeneration()); + assertEquals(context + ": version", expected.getVersion(), actual.getVersion()); + assertEquals(context + ": segment count", expected.getSegments().size(), actual.getSegments().size()); + for (int i = 0; i < expected.getSegments().size(); i++) { + Segment expectedSeg = expected.getSegments().get(i); + Segment actualSeg = actual.getSegments().get(i); + assertEquals(context + ": segment[" + i + "].generation", expectedSeg.generation(), actualSeg.generation()); + assertEquals( + context + ": segment[" + i + "].formats", + expectedSeg.dfGroupedSearchableFiles().keySet(), + actualSeg.dfGroupedSearchableFiles().keySet() + ); + for (String format : expectedSeg.dfGroupedSearchableFiles().keySet()) { + WriterFileSet expectedWfs = expectedSeg.dfGroupedSearchableFiles().get(format); + WriterFileSet actualWfs = actualSeg.dfGroupedSearchableFiles().get(format); + assertEquals(context + ": writerGeneration", expectedWfs.writerGeneration(), actualWfs.writerGeneration()); + assertEquals(context + ": files", expectedWfs.files(), actualWfs.files()); + assertEquals(context + ": numRows", expectedWfs.numRows(), actualWfs.numRows()); + } + } + assertEquals(context + ": lastWriterGeneration", expected.getLastWriterGeneration(), actual.getLastWriterGeneration()); + assertEquals(context + ": userData", expected.getUserData(), actual.getUserData()); + } + + private String generateInvalidInput(int iter) { + switch (iter % 6) { + case 0: + return randomAlphaOfLengthBetween(1, 200); + case 1: + byte[] randomBytes = new byte[randomIntBetween(1, 100)]; + random().nextBytes(randomBytes); + return java.util.Base64.getEncoder().encodeToString(randomBytes); + case 2: + DataformatAwareCatalogSnapshot snap = randomSnapshot(); + try { + String validBase64 = snap.serializeToString(); + return validBase64.substring(0, randomIntBetween(1, Math.max(1, validBase64.length() / 2))); + } catch (IOException e) { + return "AAAA"; + } + case 3: + return ""; + case 4: + return randomFrom("not-base64!!!", "===", "@@@@", "hello world", "{\"json\":true}"); + case 5: + return randomFrom("null", "undefined", "None", "nil", "NaN"); + default: + return randomAlphaOfLength(10); + } + } +} diff --git a/server/src/test/java/org/opensearch/index/engine/exec/coord/SegmentInfosCatalogSnapshotTests.java b/server/src/test/java/org/opensearch/index/engine/exec/coord/SegmentInfosCatalogSnapshotTests.java new file mode 100644 index 0000000000000..0c8986d3a1043 --- /dev/null +++ b/server/src/test/java/org/opensearch/index/engine/exec/coord/SegmentInfosCatalogSnapshotTests.java @@ -0,0 +1,87 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +package org.opensearch.index.engine.exec.coord; + +import org.apache.lucene.index.SegmentInfos; +import org.apache.lucene.util.Version; +import org.opensearch.core.common.io.stream.NamedWriteableRegistry; +import org.opensearch.test.OpenSearchTestCase; + +import java.util.Collections; +import java.util.HashMap; +import java.util.Map; + +/** + * Tests for {@link SegmentInfosCatalogSnapshot}. + */ +public class SegmentInfosCatalogSnapshotTests extends OpenSearchTestCase { + + public void testDelegation() { + for (int iter = 0; iter < 100; iter++) { + SegmentInfos segmentInfos = randomSegmentInfos(); + SegmentInfosCatalogSnapshot snapshot = new SegmentInfosCatalogSnapshot(segmentInfos); + + assertEquals(segmentInfos.getGeneration(), snapshot.getId()); + assertEquals(segmentInfos.getGeneration(), snapshot.getGeneration()); + assertEquals(segmentInfos.getVersion(), snapshot.getVersion()); + assertEquals(segmentInfos.getUserData(), snapshot.getUserData()); + assertEquals(-1L, snapshot.getLastWriterGeneration()); + assertSame(segmentInfos, snapshot.getSegmentInfos()); + + expectThrows(UnsupportedOperationException.class, snapshot::getSegments); + expectThrows(UnsupportedOperationException.class, () -> snapshot.getSearchableFiles(randomAlphaOfLength(5))); + expectThrows(UnsupportedOperationException.class, snapshot::getDataFormats); + expectThrows(UnsupportedOperationException.class, snapshot::serializeToString); + + // setUserData is a no-op, should not throw + snapshot.setUserData(Map.of("key", "value")); + } + } + + public void testClone() { + for (int iter = 0; iter < 100; iter++) { + SegmentInfos segmentInfos = randomSegmentInfos(); + SegmentInfosCatalogSnapshot snapshot = new SegmentInfosCatalogSnapshot(segmentInfos); + + SegmentInfosCatalogSnapshot cloned = snapshot.clone(); + assertNotSame(snapshot, cloned); + assertSame(segmentInfos, cloned.getSegmentInfos()); + assertEquals(snapshot.getGeneration(), cloned.getGeneration()); + assertEquals(snapshot.getVersion(), cloned.getVersion()); + } + } + + public void testCopyWriteable() throws Exception { + SegmentInfos segmentInfos = randomSegmentInfos(); + SegmentInfosCatalogSnapshot original = new SegmentInfosCatalogSnapshot(segmentInfos); + + SegmentInfosCatalogSnapshot copy = copyWriteable( + original, + new NamedWriteableRegistry(Collections.emptyList()), + SegmentInfosCatalogSnapshot::new + ); + + assertEquals(original.getGeneration(), copy.getGeneration()); + assertEquals(original.getVersion(), copy.getVersion()); + assertEquals(original.getUserData(), copy.getUserData()); + } + + // --- helpers --- + + private SegmentInfos randomSegmentInfos() { + SegmentInfos segmentInfos = new SegmentInfos(Version.LATEST.major); + int userDataEntries = randomIntBetween(0, 5); + Map userData = new HashMap<>(); + for (int i = 0; i < userDataEntries; i++) { + userData.put(randomAlphaOfLength(5), randomAlphaOfLength(10)); + } + segmentInfos.setUserData(userData, false); + return segmentInfos; + } +}