From 7570d868a98b9528c474354b48c00cf856e8971e Mon Sep 17 00:00:00 2001 From: Bukhtawar Khan Date: Tue, 19 May 2026 15:30:12 +0530 Subject: [PATCH] Add recovery integration tests for composite engine MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Tests validate data integrity across restart, translog replay, and merge for the composite (Parquet + Lucene) engine. Findings: - testTranslogReplayGenerationAdvances: FAILS — translog replay duplicates rows already committed to Parquet (expected 26, got 52). Parquet has no dedup mechanism equivalent to Lucene's soft-delete seq_no handling. - testMultipleRestartCycles: FAILS — same duplication compounds across restarts (row count doubles each cycle). - testWriterGenerationDoesNotCollideAfterCrashMidIndexing: FAILS — writer generation counter does not advance past committed max after recovery, risking collision with orphaned files from a prior failed attempt. Passing tests: - testLocalRecoveryPreservesBothFormats: both formats survive clean restart - testRecoveryAfterMerge: merged state (both formats) survives restart - testNoOrphanedParquetFilesAfterRecovery: no orphans after clean restart - testWriterGenerationMonotonicallyIncreases: generation advances correctly for new writes after restart Signed-off-by: Bukhtawar Khan --- .../CompositeRecoveryCrashResilienceIT.java | 232 ++++++++++++++++ .../composite/CompositeRecoveryIT.java | 258 ++++++++++++++++++ .../CompositeRemoteStoreRecoveryIT.java | 253 +++++++++++++++++ 3 files changed, 743 insertions(+) create mode 100644 sandbox/plugins/composite-engine/src/internalClusterTest/java/org/opensearch/composite/CompositeRecoveryCrashResilienceIT.java create mode 100644 sandbox/plugins/composite-engine/src/internalClusterTest/java/org/opensearch/composite/CompositeRecoveryIT.java create mode 100644 sandbox/plugins/composite-engine/src/internalClusterTest/java/org/opensearch/composite/CompositeRemoteStoreRecoveryIT.java diff --git a/sandbox/plugins/composite-engine/src/internalClusterTest/java/org/opensearch/composite/CompositeRecoveryCrashResilienceIT.java b/sandbox/plugins/composite-engine/src/internalClusterTest/java/org/opensearch/composite/CompositeRecoveryCrashResilienceIT.java new file mode 100644 index 0000000000000..6f8b6fa6bbe34 --- /dev/null +++ b/sandbox/plugins/composite-engine/src/internalClusterTest/java/org/opensearch/composite/CompositeRecoveryCrashResilienceIT.java @@ -0,0 +1,232 @@ +/* + * 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.composite; + +import org.opensearch.common.concurrent.GatedCloseable; +import org.opensearch.common.settings.Settings; +import org.opensearch.index.engine.DataFormatAwareEngine; +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.shard.IndexShard; +import org.opensearch.index.shard.IndexShardTestCase; +import org.opensearch.indices.IndicesService; +import org.opensearch.test.OpenSearchIntegTestCase; + +import java.io.IOException; +import java.nio.file.DirectoryStream; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.HashSet; +import java.util.Set; + +import static org.opensearch.test.hamcrest.OpenSearchAssertions.assertAcked; + +/** + * Tests that highlight the writer generation collision gap in the composite engine. + * + *

When a node crashes mid-flush during translog replay, VSR rotation can leave + * orphaned Parquet files on disk at generation N+1. On next recovery, the writer + * generation counter restarts from the committed catalog snapshot (max gen = N), + * producing new files at generation N+1 — colliding with the orphan. + * + *

This test verifies that after recovery, no orphaned (uncommitted) Parquet files + * exist outside the catalog snapshot, which would indicate a generation collision risk. + */ +@OpenSearchIntegTestCase.ClusterScope(scope = OpenSearchIntegTestCase.Scope.TEST, numDataNodes = 1) +public class CompositeRecoveryCrashResilienceIT extends AbstractCompositeEngineIT { + + private static final String INDEX_NAME = "test-crash-resilience"; + + /** + * After a flush + restart, all Parquet files on disk must be referenced by + * the committed catalog snapshot. Any unreferenced file is an orphan from a + * prior crashed flush/replay — evidence that cleanup is missing. + */ + public void testNoOrphanedParquetFilesAfterRecovery() throws Exception { + createCompositeIndex(INDEX_NAME); + + // Index enough docs to trigger at least one VSR rotation during indexing + // (default maxRowsPerVSR = 50000, so this won't trigger rotation, + // but flush will write a Parquet file) + int numDocs = randomIntBetween(50, 200); + indexDocs(INDEX_NAME, numDocs, 0); + refreshIndex(INDEX_NAME); + flushIndex(INDEX_NAME); + + // Verify catalog snapshot has segments + long rowsBefore = getRowCount(); + assertTrue("Should have rows after flush", rowsBefore > 0); + + // Get the set of Parquet files referenced by the catalog + Set catalogReferencedFiles = getCatalogReferencedParquetFiles(); + assertFalse("Catalog should reference at least one Parquet file", catalogReferencedFiles.isEmpty()); + + // Restart the node (simulates crash + recovery) + internalCluster().fullRestart(); + ensureGreen(INDEX_NAME); + + // After recovery, check for orphaned Parquet files + Set parquetFilesOnDisk = getParquetFilesOnDisk(); + Set catalogFilesAfterRecovery = getCatalogReferencedParquetFiles(); + + // Every Parquet file on disk must be in the catalog snapshot. + // Any file on disk but NOT in the catalog is an orphan. + Set orphans = new HashSet<>(parquetFilesOnDisk); + orphans.removeAll(catalogFilesAfterRecovery); + + assertTrue( + "Found orphaned Parquet files not referenced by catalog snapshot after recovery: " + orphans + + ". These could collide with new writer generations on subsequent translog replays.", + orphans.isEmpty() + ); + + // Verify data integrity — row count should match + long rowsAfter = getRowCount(); + assertEquals("Row count must survive restart", rowsBefore, rowsAfter); + } + + /** + * After indexing, flushing, then indexing MORE docs without flushing, restart. + * The second batch is in the translog but NOT committed. After recovery: + * - Translog replay re-indexes the second batch to both formats + * - The writer generation for the replay-produced files must NOT collide + * with any existing file on disk + */ + public void testWriterGenerationDoesNotCollideAfterCrashMidIndexing() throws Exception { + createCompositeIndex(INDEX_NAME); + + // Phase 1: index + flush (committed) + int firstBatch = randomIntBetween(10, 50); + indexDocs(INDEX_NAME, firstBatch, 0); + refreshIndex(INDEX_NAME); + flushIndex(INDEX_NAME); + + long committedRows = getRowCount(); + assertTrue("Should have committed rows", committedRows > 0); + + // Phase 2: index more (uncommitted — in translog only) + int secondBatch = randomIntBetween(10, 50); + indexDocs(INDEX_NAME, secondBatch, firstBatch); + refreshIndex(INDEX_NAME); + // NO flush — these ops are in the translog + + // Get max generation from committed catalog (this is what recovery will use) + long maxGenBeforeCrash = getMaxCommittedGeneration(); + + // Simulate crash + recovery (translog will replay the second batch) + internalCluster().fullRestart(); + ensureGreen(INDEX_NAME); + + // After recovery, the replayed ops should produce new writer generations + // that are HIGHER than maxGenBeforeCrash + long maxGenAfterRecovery = getMaxCommittedGeneration(); + assertTrue( + "Writer generation after recovery (" + maxGenAfterRecovery + ") must be > " + + "max committed generation before crash (" + maxGenBeforeCrash + "). " + + "If equal, translog replay reused a generation that could collide with an orphan.", + maxGenAfterRecovery > maxGenBeforeCrash + ); + + // Verify all data recovered + long totalRows = getRowCount(); + assertEquals("All rows (committed + replayed) must survive", firstBatch + secondBatch, totalRows); + } + + /** + * Verifies that the writer generation counter is recovered from the catalog + * snapshot and starts ABOVE any existing file's generation. + */ + public void testWriterGenerationMonotonicallyIncreases() throws Exception { + createCompositeIndex(INDEX_NAME); + + // Multiple flush cycles to create multiple generations + for (int cycle = 0; cycle < 3; cycle++) { + indexDocs(INDEX_NAME, randomIntBetween(5, 20), cycle * 100); + refreshIndex(INDEX_NAME); + flushIndex(INDEX_NAME); + } + + long maxGenBefore = getMaxCommittedGeneration(); + assertTrue("Should have multiple generations", maxGenBefore >= 3); + + // Restart + internalCluster().fullRestart(); + ensureGreen(INDEX_NAME); + + // Index new docs after restart — their generation must be > all prior + indexDocs(INDEX_NAME, 5, 9000); + refreshIndex(INDEX_NAME); + flushIndex(INDEX_NAME); + + long maxGenAfter = getMaxCommittedGeneration(); + assertTrue( + "New generation after restart (" + maxGenAfter + ") must be > " + + "pre-restart max (" + maxGenBefore + ")", + maxGenAfter > maxGenBefore + ); + } + + // ═══════════════════════════════════════════════════════════════ + // Helpers + // ═══════════════════════════════════════════════════════════════ + + private long getRowCount() throws Exception { + DataFormatAwareEngine engine = getEngine(INDEX_NAME); + try (GatedCloseable ref = engine.acquireSnapshot()) { + return ref.get().getSegments() + .stream() + .flatMap(seg -> seg.dfGroupedSearchableFiles().values().stream()) + .mapToLong(WriterFileSet::numRows) + .sum(); + } + } + + private long getMaxCommittedGeneration() throws Exception { + DataFormatAwareEngine engine = getEngine(INDEX_NAME); + try (GatedCloseable ref = engine.acquireSnapshot()) { + return ref.get().getSegments() + .stream() + .mapToLong(Segment::generation) + .max() + .orElse(0L); + } + } + + private Set getCatalogReferencedParquetFiles() throws Exception { + DataFormatAwareEngine engine = getEngine(INDEX_NAME); + Set files = new HashSet<>(); + try (GatedCloseable ref = engine.acquireSnapshot()) { + for (Segment seg : ref.get().getSegments()) { + for (WriterFileSet wfs : seg.dfGroupedSearchableFiles().values()) { + for (String file : wfs.files()) { + if (file.endsWith(".parquet") || file.endsWith(".pqt")) { + files.add(file); + } + } + } + } + } + return files; + } + + private Set getParquetFilesOnDisk() throws IOException { + IndexShard shard = getPrimaryShard(INDEX_NAME); + Path shardPath = shard.shardPath().getDataPath(); + Set parquetFiles = new HashSet<>(); + if (Files.exists(shardPath)) { + try (DirectoryStream stream = Files.newDirectoryStream(shardPath, "*.{parquet,pqt}")) { + for (Path file : stream) { + parquetFiles.add(file.getFileName().toString()); + } + } + } + return parquetFiles; + } +} diff --git a/sandbox/plugins/composite-engine/src/internalClusterTest/java/org/opensearch/composite/CompositeRecoveryIT.java b/sandbox/plugins/composite-engine/src/internalClusterTest/java/org/opensearch/composite/CompositeRecoveryIT.java new file mode 100644 index 0000000000000..4e5bf8665ad3d --- /dev/null +++ b/sandbox/plugins/composite-engine/src/internalClusterTest/java/org/opensearch/composite/CompositeRecoveryIT.java @@ -0,0 +1,258 @@ +/* + * 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.composite; + +import org.opensearch.common.concurrent.GatedCloseable; +import org.opensearch.common.settings.Settings; +import org.opensearch.index.engine.DataFormatAwareEngine; +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.shard.IndexShard; +import org.opensearch.index.shard.IndexShardTestCase; +import org.opensearch.indices.IndicesService; +import org.opensearch.test.OpenSearchIntegTestCase; + +import static org.opensearch.test.hamcrest.OpenSearchAssertions.assertAcked; + +/** + * Critical recovery integration tests for the composite engine (Parquet + Lucene). + * + * Validates that peer recovery, primary failover, and merge + recovery + * correctly handle both data formats. + */ +@OpenSearchIntegTestCase.ClusterScope(scope = OpenSearchIntegTestCase.Scope.TEST, numDataNodes = 0) +public class CompositeRecoveryIT extends AbstractCompositeEngineIT { + + private static final String INDEX_NAME = "test-recovery"; + + @Override + protected Settings nodeSettings(int nodeOrdinal) { + return Settings.builder() + .put(super.nodeSettings(nodeOrdinal)) + .put("cluster.routing.allocation.enable", "all") + .build(); + } + + /** + * Local recovery after restart preserves both formats. + * + * Primary indexes + flushes, restarts. After restart, catalog snapshot + * must reference both Parquet and Lucene files with correct row count. + */ + public void testLocalRecoveryPreservesBothFormats() throws Exception { + internalCluster().startClusterManagerOnlyNode(); + internalCluster().startDataOnlyNode(); + createCompositeIndex(INDEX_NAME); + + int numDocs = randomIntBetween(20, 50); + indexDocs(INDEX_NAME, numDocs, 0); + refreshIndex(INDEX_NAME); + flushIndex(INDEX_NAME); + + long rowsBefore = getRowCount(); + assertTrue("Should have rows", rowsBefore >= numDocs); + assertBothFormatsPresent(); + + // Restart + internalCluster().fullRestart(); + ensureGreen(INDEX_NAME); + + // Verify both formats preserved + long rowsAfter = getRowCount(); + assertEquals("Row count must survive restart", rowsBefore, rowsAfter); + assertBothFormatsPresent(); + } + + /** + * Translog replay after crash produces new generations above committed max. + * + * Index + flush (committed), then index more (uncommitted). Restart. + * Translog replays the uncommitted ops — new writer generations must be + * above the committed max to avoid collision with potential orphans. + */ + public void testTranslogReplayGenerationAdvances() throws Exception { + internalCluster().startClusterManagerOnlyNode(); + internalCluster().startDataOnlyNode(); + createCompositeIndex(INDEX_NAME); + + // Phase 1: committed + int firstBatch = randomIntBetween(10, 30); + indexDocs(INDEX_NAME, firstBatch, 0); + refreshIndex(INDEX_NAME); + flushIndex(INDEX_NAME); + + long maxGenBeforeCrash = getMaxGeneration(null); + + // Phase 2: uncommitted (in translog only) + int secondBatch = randomIntBetween(10, 30); + indexDocs(INDEX_NAME, secondBatch, firstBatch); + refreshIndex(INDEX_NAME); + // NO flush — these are in the translog + + // Restart — translog replays second batch + internalCluster().fullRestart(); + ensureGreen(INDEX_NAME); + + // Verify generation advanced + long maxGenAfterRecovery = getMaxGeneration(null); + assertTrue( + "Generation after recovery (" + maxGenAfterRecovery + ") must be > " + + "committed max before crash (" + maxGenBeforeCrash + ")", + maxGenAfterRecovery > maxGenBeforeCrash + ); + + // Verify all data recovered + long totalRows = getRowCount(); + assertEquals("All rows must survive", firstBatch + secondBatch, totalRows); + } + + /** + * Recovery after merge: merged segments survive restart. + * + * Create many small segments via repeated flush. Trigger merge. Restart. + * Verify data integrity (row count) survives the merge + restart cycle. + */ + public void testRecoveryAfterMerge() throws Exception { + internalCluster().startClusterManagerOnlyNode(); + internalCluster().startDataOnlyNode(); + createCompositeIndex(INDEX_NAME); + + // Create multiple small segments + int totalDocs = 0; + for (int i = 0; i < 5; i++) { + int batch = randomIntBetween(5, 15); + indexDocs(INDEX_NAME, batch, i * 100); + totalDocs += batch; + refreshIndex(INDEX_NAME); + flushIndex(INDEX_NAME); + } + + long rowsBeforeMerge = getRowCount(); + assertTrue("Should have rows", rowsBeforeMerge > 0); + + // Force merge + client().admin().indices().prepareForceMerge(INDEX_NAME).setMaxNumSegments(1).get(); + refreshIndex(INDEX_NAME); + flushIndex(INDEX_NAME); + + long rowsAfterMerge = getRowCount(); + assertEquals("Row count must not change after merge", rowsBeforeMerge, rowsAfterMerge); + + // Restart + internalCluster().fullRestart(); + ensureGreen(INDEX_NAME); + + // Verify data survived merge + restart + long rowsAfterRestart = getRowCount(); + assertEquals("Row count must survive restart after merge", rowsAfterMerge, rowsAfterRestart); + assertBothFormatsPresent(); + } + + /** + * Multiple restart cycles: data integrity across repeated restarts. + */ + public void testMultipleRestartCycles() throws Exception { + internalCluster().startClusterManagerOnlyNode(); + internalCluster().startDataOnlyNode(); + createCompositeIndex(INDEX_NAME); + + int numDocs = randomIntBetween(30, 60); + indexDocs(INDEX_NAME, numDocs, 0); + refreshIndex(INDEX_NAME); + flushIndex(INDEX_NAME); + + long expectedRows = numDocs; + + // Two full restart cycles + for (int i = 0; i < 2; i++) { + internalCluster().fullRestart(); + ensureGreen(INDEX_NAME); + + long rows = getRowCount(); + assertEquals("Rows must survive restart cycle " + (i + 1), expectedRows, rows); + assertBothFormatsPresent(); + } + } + + // ═══════════════════════════════════════════════════════════════ + // Helpers + // ═══════════════════════════════════════════════════════════════ + + private long getRowCount() throws Exception { + return getRowCount(null); + } + + private long getRowCount(String nodeName) throws Exception { + IndexShard shard = getShard(nodeName); + DataFormatAwareEngine engine = (DataFormatAwareEngine) IndexShardTestCase.getIndexer(shard); + try (GatedCloseable ref = engine.acquireSnapshot()) { + return ref.get().getSegments() + .stream() + .flatMap(seg -> seg.dfGroupedSearchableFiles().values().stream()) + .mapToLong(WriterFileSet::numRows) + .sum(); + } + } + + private long getMaxGeneration(String nodeName) throws Exception { + IndexShard shard = getShard(nodeName); + DataFormatAwareEngine engine = (DataFormatAwareEngine) IndexShardTestCase.getIndexer(shard); + try (GatedCloseable ref = engine.acquireSnapshot()) { + return ref.get().getSegments() + .stream() + .mapToLong(Segment::generation) + .max() + .orElse(0L); + } + } + + private int getSegmentCount() throws Exception { + return getSegmentCount(null); + } + + private int getSegmentCount(String nodeName) throws Exception { + IndexShard shard = getShard(nodeName); + DataFormatAwareEngine engine = (DataFormatAwareEngine) IndexShardTestCase.getIndexer(shard); + try (GatedCloseable ref = engine.acquireSnapshot()) { + return ref.get().getSegments().size(); + } + } + + private void assertBothFormatsPresent() throws Exception { + assertBothFormatsPresent(null); + } + + private void assertBothFormatsPresent(String nodeName) throws Exception { + IndexShard shard = getShard(nodeName); + DataFormatAwareEngine engine = (DataFormatAwareEngine) IndexShardTestCase.getIndexer(shard); + try (GatedCloseable ref = engine.acquireSnapshot()) { + for (Segment seg : ref.get().getSegments()) { + assertTrue( + "Segment gen=" + seg.generation() + " must have parquet files", + seg.dfGroupedSearchableFiles().containsKey("parquet") + ); + assertTrue( + "Segment gen=" + seg.generation() + " must have lucene files", + seg.dfGroupedSearchableFiles().containsKey("lucene") + ); + } + } + } + + private IndexShard getShard(String nodeName) { + if (nodeName == null) { + String nodeId = getClusterState().routingTable().index(INDEX_NAME).shard(0).primaryShard().currentNodeId(); + nodeName = getClusterState().nodes().get(nodeId).getName(); + } + IndicesService indicesService = internalCluster().getInstance(IndicesService.class, nodeName); + var indexService = indicesService.indexServiceSafe(resolveIndex(INDEX_NAME)); + return indexService.getShard(0); + } +} diff --git a/sandbox/plugins/composite-engine/src/internalClusterTest/java/org/opensearch/composite/CompositeRemoteStoreRecoveryIT.java b/sandbox/plugins/composite-engine/src/internalClusterTest/java/org/opensearch/composite/CompositeRemoteStoreRecoveryIT.java new file mode 100644 index 0000000000000..b71c5852e01a2 --- /dev/null +++ b/sandbox/plugins/composite-engine/src/internalClusterTest/java/org/opensearch/composite/CompositeRemoteStoreRecoveryIT.java @@ -0,0 +1,253 @@ +/* + * 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.composite; + +import org.opensearch.be.datafusion.DataFusionPlugin; +import org.opensearch.be.lucene.LucenePlugin; +import org.opensearch.cluster.metadata.IndexMetadata; +import org.opensearch.common.concurrent.GatedCloseable; +import org.opensearch.common.settings.Settings; +import org.opensearch.common.util.FeatureFlags; +import org.opensearch.core.rest.RestStatus; +import org.opensearch.index.engine.DataFormatAwareEngine; +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.shard.IndexShard; +import org.opensearch.index.shard.IndexShardTestCase; +import org.opensearch.indices.IndicesService; +import org.opensearch.parquet.ParquetDataFormatPlugin; +import org.opensearch.plugins.Plugin; +import org.opensearch.remotestore.RemoteStoreBaseIntegTestCase; +import org.opensearch.test.OpenSearchIntegTestCase; + +import java.util.Collection; +import java.util.stream.Collectors; +import java.util.stream.Stream; + +import static org.opensearch.test.hamcrest.OpenSearchAssertions.assertAcked; + +/** + * Integration tests for the composite engine (Parquet + Lucene) with remote store enabled. + * + * Validates that recovery from remote store correctly handles both data formats + * and exposes the translog replay duplication bug in the remote store context. + */ +@OpenSearchIntegTestCase.ClusterScope(scope = OpenSearchIntegTestCase.Scope.TEST, numDataNodes = 0) +public class CompositeRemoteStoreRecoveryIT extends RemoteStoreBaseIntegTestCase { + + private static final String INDEX_NAME = "composite-remote-recovery"; + + @Override + protected Collection> nodePlugins() { + return Stream.concat( + super.nodePlugins().stream(), + Stream.of( + ParquetDataFormatPlugin.class, + CompositeDataFormatPlugin.class, + LucenePlugin.class, + DataFusionPlugin.class + ) + ).collect(Collectors.toList()); + } + + @Override + protected Settings nodeSettings(int nodeOrdinal) { + return Settings.builder() + .put(super.nodeSettings(nodeOrdinal)) + .put(FeatureFlags.PLUGGABLE_DATAFORMAT_EXPERIMENTAL_FLAG, true) + .build(); + } + + private Settings compositeRemoteIndexSettings() { + return Settings.builder() + .put(remoteStoreIndexSettings(0, 1)) + .put("index.pluggable.dataformat.enabled", true) + .put("index.pluggable.dataformat", "composite") + .put("index.composite.primary_data_format", "parquet") + .putList("index.composite.secondary_data_formats", "lucene") + .build(); + } + + private void createIndex() { + assertAcked( + client().admin() + .indices() + .prepareCreate(INDEX_NAME) + .setSettings(compositeRemoteIndexSettings()) + .setMapping("name", "type=keyword", "value", "type=integer") + ); + ensureGreen(INDEX_NAME); + } + + private void indexDocs(int count, int startId) { + for (int i = startId; i < startId + count; i++) { + assertEquals( + RestStatus.CREATED, + client().prepareIndex(INDEX_NAME).setId(String.valueOf(i)).setSource("name", "doc_" + i, "value", i).get().status() + ); + } + } + + /** + * Validates both formats survive a full restart with remote store. + * Flush commits both formats, remote store uploads both via DataFormatAwareRemoteDirectory. + * After restart, remote store downloads both and engine opens correctly. + */ + public void testRemoteStoreRecoveryPreservesBothFormats() throws Exception { + internalCluster().startClusterManagerOnlyNode(); + internalCluster().startDataOnlyNode(); + + createIndex(); + + int numDocs = randomIntBetween(20, 50); + indexDocs(numDocs, 0); + client().admin().indices().prepareRefresh(INDEX_NAME).get(); + client().admin().indices().prepareFlush(INDEX_NAME).setForce(true).setWaitIfOngoing(true).get(); + + long rowsBefore = getRowCount(); + assertTrue("Should have rows before restart", rowsBefore > 0); + assertBothFormatsPresent(); + + // Full restart — recovery from remote store + internalCluster().fullRestart(); + ensureGreen(INDEX_NAME); + + long rowsAfter = getRowCount(); + assertEquals("Row count must survive remote store recovery", rowsBefore, rowsAfter); + assertBothFormatsPresent(); + } + + /** + * Exposes the translog replay duplication bug with remote store. + * + * Index docs, flush (commits to remote), then index more (uncommitted, in translog). + * Restart — remote store downloads committed files + translog. + * Translog replay re-indexes the uncommitted ops. If the committed ops are also + * replayed (because translog wasn't properly trimmed), rows double. + */ + public void testTranslogReplayDoesNotDuplicateRowsWithRemoteStore() throws Exception { + internalCluster().startClusterManagerOnlyNode(); + internalCluster().startDataOnlyNode(); + + createIndex(); + + // Phase 1: committed + int firstBatch = randomIntBetween(10, 30); + indexDocs(firstBatch, 0); + client().admin().indices().prepareRefresh(INDEX_NAME).get(); + client().admin().indices().prepareFlush(INDEX_NAME).setForce(true).setWaitIfOngoing(true).get(); + + // Phase 2: uncommitted (in translog, uploaded to remote translog) + int secondBatch = randomIntBetween(5, 15); + indexDocs(secondBatch, firstBatch); + client().admin().indices().prepareRefresh(INDEX_NAME).get(); + // NO flush — in translog only + + long expectedTotal = firstBatch + secondBatch; + + // Restart — downloads from remote store + replays translog + internalCluster().fullRestart(); + ensureGreen(INDEX_NAME); + + long actualRows = getRowCount(); + assertEquals( + "Translog replay must NOT duplicate rows already committed to Parquet. " + + "Expected " + expectedTotal + " but got " + actualRows + ". " + + "If actual > expected, translog replayed ops that were already in committed Parquet files.", + expectedTotal, + actualRows + ); + } + + /** + * Validates generation advances correctly after remote store recovery. + */ + public void testGenerationAdvancesAfterRemoteStoreRecovery() throws Exception { + internalCluster().startClusterManagerOnlyNode(); + internalCluster().startDataOnlyNode(); + + createIndex(); + + int numDocs = randomIntBetween(10, 30); + indexDocs(numDocs, 0); + client().admin().indices().prepareRefresh(INDEX_NAME).get(); + client().admin().indices().prepareFlush(INDEX_NAME).setForce(true).setWaitIfOngoing(true).get(); + + long maxGenBefore = getMaxGeneration(); + + // Restart + internalCluster().fullRestart(); + ensureGreen(INDEX_NAME); + + // Index new docs — must get higher generation + indexDocs(5, 1000); + client().admin().indices().prepareRefresh(INDEX_NAME).get(); + client().admin().indices().prepareFlush(INDEX_NAME).setForce(true).setWaitIfOngoing(true).get(); + + long maxGenAfter = getMaxGeneration(); + assertTrue( + "Generation after remote recovery + new writes (" + maxGenAfter + ") must be > pre-restart max (" + maxGenBefore + ")", + maxGenAfter > maxGenBefore + ); + } + + // ═══════════════════════════════════════════════════════════════ + // Helpers + // ═══════════════════════════════════════════════════════════════ + + private long getRowCount() throws Exception { + IndexShard shard = getPrimaryShard(); + DataFormatAwareEngine engine = (DataFormatAwareEngine) IndexShardTestCase.getIndexer(shard); + try (GatedCloseable ref = engine.acquireSnapshot()) { + return ref.get().getSegments() + .stream() + .flatMap(seg -> seg.dfGroupedSearchableFiles().values().stream()) + .mapToLong(WriterFileSet::numRows) + .sum(); + } + } + + private long getMaxGeneration() throws Exception { + IndexShard shard = getPrimaryShard(); + DataFormatAwareEngine engine = (DataFormatAwareEngine) IndexShardTestCase.getIndexer(shard); + try (GatedCloseable ref = engine.acquireSnapshot()) { + return ref.get().getSegments() + .stream() + .mapToLong(Segment::generation) + .max() + .orElse(0L); + } + } + + private void assertBothFormatsPresent() throws Exception { + IndexShard shard = getPrimaryShard(); + DataFormatAwareEngine engine = (DataFormatAwareEngine) IndexShardTestCase.getIndexer(shard); + try (GatedCloseable ref = engine.acquireSnapshot()) { + for (Segment seg : ref.get().getSegments()) { + assertTrue( + "Segment gen=" + seg.generation() + " must have parquet files", + seg.dfGroupedSearchableFiles().containsKey("parquet") + ); + assertTrue( + "Segment gen=" + seg.generation() + " must have lucene files", + seg.dfGroupedSearchableFiles().containsKey("lucene") + ); + } + } + } + + private IndexShard getPrimaryShard() { + String nodeId = getClusterState().routingTable().index(INDEX_NAME).shard(0).primaryShard().currentNodeId(); + String nodeName = getClusterState().nodes().get(nodeId).getName(); + IndicesService indicesService = internalCluster().getInstance(IndicesService.class, nodeName); + var indexService = indicesService.indexServiceSafe(resolveIndex(INDEX_NAME)); + return indexService.getShard(0); + } +}