diff --git a/sandbox/plugins/composite-engine/build.gradle b/sandbox/plugins/composite-engine/build.gradle index ba7c3a12f0b98..bfd41f508ead8 100644 --- a/sandbox/plugins/composite-engine/build.gradle +++ b/sandbox/plugins/composite-engine/build.gradle @@ -38,7 +38,6 @@ internalClusterTest { dependencies { api project(':libs:opensearch-concurrent-queue') - api project(':sandbox:libs:composite-common') compileOnly project(':server') testImplementation project(':test:framework') testImplementation project(':sandbox:plugins:parquet-data-format') diff --git a/sandbox/plugins/composite-engine/src/internalClusterTest/java/org/opensearch/composite/CompositeMergeIT.java b/sandbox/plugins/composite-engine/src/internalClusterTest/java/org/opensearch/composite/CompositeMergeIT.java index 9be2513ef2d34..0ce2f195a4b0d 100644 --- a/sandbox/plugins/composite-engine/src/internalClusterTest/java/org/opensearch/composite/CompositeMergeIT.java +++ b/sandbox/plugins/composite-engine/src/internalClusterTest/java/org/opensearch/composite/CompositeMergeIT.java @@ -8,6 +8,8 @@ package org.opensearch.composite; +import com.carrotsearch.randomizedtesting.annotations.ThreadLeakScope; + import org.opensearch.action.admin.indices.refresh.RefreshResponse; import org.opensearch.action.admin.indices.stats.IndicesStatsResponse; import org.opensearch.action.admin.indices.stats.ShardStats; @@ -48,13 +50,9 @@ import java.util.Set; import java.util.function.Function; -/** - * Integration tests for composite merge with real Parquet backend. - * - * Run with: - * ./gradlew :sandbox:plugins:composite-engine:internalClusterTest \ - * --tests "*.CompositeMergeIT" -Dsandbox.enabled=true - */ +// The Tokio IO runtime worker thread (used by the Rust merge k-way merge sort) is a process-lifetime +// singleton that persists after tests complete. It polls for new async IO tasks between merges. +@ThreadLeakScope(ThreadLeakScope.Scope.NONE) @OpenSearchIntegTestCase.ClusterScope(scope = OpenSearchIntegTestCase.Scope.TEST, numDataNodes = 1) public class CompositeMergeIT extends OpenSearchIntegTestCase { diff --git a/sandbox/plugins/parquet-data-format/benchmarks/build.gradle b/sandbox/plugins/parquet-data-format/benchmarks/build.gradle index ee90cb6d2301b..137d589e558cd 100644 --- a/sandbox/plugins/parquet-data-format/benchmarks/build.gradle +++ b/sandbox/plugins/parquet-data-format/benchmarks/build.gradle @@ -54,7 +54,7 @@ dependencies { api "org.slf4j:slf4j-api:${versions.slf4j}" api "org.apache.logging.log4j:log4j-api:${versions.log4j}" api "org.apache.logging.log4j:log4j-core:${versions.log4j}" - api "org.apache.logging.log4j:log4j-slf4j-impl:${versions.log4j}" + api "org.apache.logging.log4j:log4j-slf4j2-impl:${versions.log4j}" } // enable the JMH's BenchmarkProcessor to generate the final benchmark classes diff --git a/sandbox/plugins/parquet-data-format/build.gradle b/sandbox/plugins/parquet-data-format/build.gradle index 3e1547568798e..8ba011b0b92fb 100644 --- a/sandbox/plugins/parquet-data-format/build.gradle +++ b/sandbox/plugins/parquet-data-format/build.gradle @@ -71,3 +71,10 @@ test { tasks.matching { it.name == 'missingJavadoc' }.configureEach { enabled = false } + +tasks.named('forbiddenPatterns').configure { + exclude '**/*.dylib' + exclude '**/*.so' + exclude '**/*.dll' + exclude '**/*.parquet' +} diff --git a/sandbox/plugins/parquet-data-format/licenses/log4j-slf4j2-impl-2.25.4.jar.sha1 b/sandbox/plugins/parquet-data-format/licenses/log4j-slf4j2-impl-2.25.4.jar.sha1 new file mode 100644 index 0000000000000..f018d071914e4 --- /dev/null +++ b/sandbox/plugins/parquet-data-format/licenses/log4j-slf4j2-impl-2.25.4.jar.sha1 @@ -0,0 +1 @@ +052a8e43b29eee3b9d6cd9bad696f5d2284d7053 \ No newline at end of file diff --git a/sandbox/plugins/parquet-data-format/licenses/log4j-slf4j2-impl-LICENSE.txt b/sandbox/plugins/parquet-data-format/licenses/log4j-slf4j2-impl-LICENSE.txt new file mode 100644 index 0000000000000..6279e5206de13 --- /dev/null +++ b/sandbox/plugins/parquet-data-format/licenses/log4j-slf4j2-impl-LICENSE.txt @@ -0,0 +1,202 @@ + + Apache License + Version 2.0, January 2004 + http://www.apache.org/licenses/ + + TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION + + 1. Definitions. + + "License" shall mean the terms and conditions for use, reproduction, + and distribution as defined by Sections 1 through 9 of this document. + + "Licensor" shall mean the copyright owner or entity authorized by + the copyright owner that is granting the License. + + "Legal Entity" shall mean the union of the acting entity and all + other entities that control, are controlled by, or are under common + control with that entity. For the purposes of this definition, + "control" means (i) the power, direct or indirect, to cause the + direction or management of such entity, whether by contract or + otherwise, or (ii) ownership of fifty percent (50%) or more of the + outstanding shares, or (iii) beneficial ownership of such entity. + + "You" (or "Your") shall mean an individual or Legal Entity + exercising permissions granted by this License. + + "Source" form shall mean the preferred form for making modifications, + including but not limited to software source code, documentation + source, and configuration files. + + "Object" form shall mean any form resulting from mechanical + transformation or translation of a Source form, including but + not limited to compiled object code, generated documentation, + and conversions to other media types. + + "Work" shall mean the work of authorship, whether in Source or + Object form, made available under the License, as indicated by a + copyright notice that is included in or attached to the work + (an example is provided in the Appendix below). + + "Derivative Works" shall mean any work, whether in Source or Object + form, that is based on (or derived from) the Work and for which the + editorial revisions, annotations, elaborations, or other modifications + represent, as a whole, an original work of authorship. For the purposes + of this License, Derivative Works shall not include works that remain + separable from, or merely link (or bind by name) to the interfaces of, + the Work and Derivative Works thereof. + + "Contribution" shall mean any work of authorship, including + the original version of the Work and any modifications or additions + to that Work or Derivative Works thereof, that is intentionally + submitted to Licensor for inclusion in the Work by the copyright owner + or by an individual or Legal Entity authorized to submit on behalf of + the copyright owner. For the purposes of this definition, "submitted" + means any form of electronic, verbal, or written communication sent + to the Licensor or its representatives, including but not limited to + communication on electronic mailing lists, source code control systems, + and issue tracking systems that are managed by, or on behalf of, the + Licensor for the purpose of discussing and improving the Work, but + excluding communication that is conspicuously marked or otherwise + designated in writing by the copyright owner as "Not a Contribution." + + "Contributor" shall mean Licensor and any individual or Legal Entity + on behalf of whom a Contribution has been received by Licensor and + subsequently incorporated within the Work. + + 2. Grant of Copyright License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + copyright license to reproduce, prepare Derivative Works of, + publicly display, publicly perform, sublicense, and distribute the + Work and such Derivative Works in Source or Object form. + + 3. Grant of Patent License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + (except as stated in this section) patent license to make, have made, + use, offer to sell, sell, import, and otherwise transfer the Work, + where such license applies only to those patent claims licensable + by such Contributor that are necessarily infringed by their + Contribution(s) alone or by combination of their Contribution(s) + with the Work to which such Contribution(s) was submitted. If You + institute patent litigation against any entity (including a + cross-claim or counterclaim in a lawsuit) alleging that the Work + or a Contribution incorporated within the Work constitutes direct + or contributory patent infringement, then any patent licenses + granted to You under this License for that Work shall terminate + as of the date such litigation is filed. + + 4. Redistribution. You may reproduce and distribute copies of the + Work or Derivative Works thereof in any medium, with or without + modifications, and in Source or Object form, provided that You + meet the following conditions: + + (a) You must give any other recipients of the Work or + Derivative Works a copy of this License; and + + (b) You must cause any modified files to carry prominent notices + stating that You changed the files; and + + (c) You must retain, in the Source form of any Derivative Works + that You distribute, all copyright, patent, trademark, and + attribution notices from the Source form of the Work, + excluding those notices that do not pertain to any part of + the Derivative Works; and + + (d) If the Work includes a "NOTICE" text file as part of its + distribution, then any Derivative Works that You distribute must + include a readable copy of the attribution notices contained + within such NOTICE file, excluding those notices that do not + pertain to any part of the Derivative Works, in at least one + of the following places: within a NOTICE text file distributed + as part of the Derivative Works; within the Source form or + documentation, if provided along with the Derivative Works; or, + within a display generated by the Derivative Works, if and + wherever such third-party notices normally appear. The contents + of the NOTICE file are for informational purposes only and + do not modify the License. You may add Your own attribution + notices within Derivative Works that You distribute, alongside + or as an addendum to the NOTICE text from the Work, provided + that such additional attribution notices cannot be construed + as modifying the License. + + You may add Your own copyright statement to Your modifications and + may provide additional or different license terms and conditions + for use, reproduction, or distribution of Your modifications, or + for any such Derivative Works as a whole, provided Your use, + reproduction, and distribution of the Work otherwise complies with + the conditions stated in this License. + + 5. Submission of Contributions. Unless You explicitly state otherwise, + any Contribution intentionally submitted for inclusion in the Work + by You to the Licensor shall be under the terms and conditions of + this License, without any additional terms or conditions. + Notwithstanding the above, nothing herein shall supersede or modify + the terms of any separate license agreement you may have executed + with Licensor regarding such Contributions. + + 6. Trademarks. This License does not grant permission to use the trade + names, trademarks, service marks, or product names of the Licensor, + except as required for reasonable and customary use in describing the + origin of the Work and reproducing the content of the NOTICE file. + + 7. Disclaimer of Warranty. Unless required by applicable law or + agreed to in writing, Licensor provides the Work (and each + Contributor provides its Contributions) on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or + implied, including, without limitation, any warranties or conditions + of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A + PARTICULAR PURPOSE. You are solely responsible for determining the + appropriateness of using or redistributing the Work and assume any + risks associated with Your exercise of permissions under this License. + + 8. Limitation of Liability. In no event and under no legal theory, + whether in tort (including negligence), contract, or otherwise, + unless required by applicable law (such as deliberate and grossly + negligent acts) or agreed to in writing, shall any Contributor be + liable to You for damages, including any direct, indirect, special, + incidental, or consequential damages of any character arising as a + result of this License or out of the use or inability to use the + Work (including but not limited to damages for loss of goodwill, + work stoppage, computer failure or malfunction, or any and all + other commercial damages or losses), even if such Contributor + has been advised of the possibility of such damages. + + 9. Accepting Warranty or Additional Liability. While redistributing + the Work or Derivative Works thereof, You may choose to offer, + and charge a fee for, acceptance of support, warranty, indemnity, + or other liability obligations and/or rights consistent with this + License. However, in accepting such obligations, You may act only + on Your own behalf and on Your sole responsibility, not on behalf + of any other Contributor, and only if You agree to indemnify, + defend, and hold each Contributor harmless for any liability + incurred by, or claims asserted against, such Contributor by reason + of your accepting any such warranty or additional liability. + + END OF TERMS AND CONDITIONS + + APPENDIX: How to apply the Apache License to your work. + + To apply the Apache License to your work, attach the following + boilerplate notice, with the fields enclosed by brackets "[]" + replaced with your own identifying information. (Don't include + the brackets!) The text should be enclosed in the appropriate + comment syntax for the file format. We also recommend that a + file or class name and description of purpose be included on the + same "printed page" as the copyright notice for easier + identification within third-party archives. + + Copyright 1999-2005 The Apache Software Foundation + + Licensed under the Apache License, Version 2.0 (the "License"); + you may not use this file except in compliance with the License. + You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. diff --git a/sandbox/plugins/parquet-data-format/licenses/log4j-slf4j2-impl-NOTICE.txt b/sandbox/plugins/parquet-data-format/licenses/log4j-slf4j2-impl-NOTICE.txt new file mode 100644 index 0000000000000..5a296bfcd19ec --- /dev/null +++ b/sandbox/plugins/parquet-data-format/licenses/log4j-slf4j2-impl-NOTICE.txt @@ -0,0 +1,6 @@ +SLF4J 2 Provider for Log4j API +Copyright 1999-2025 The Apache Software Foundation + + +This product includes software developed at +The Apache Software Foundation (http://www.apache.org/). diff --git a/sandbox/plugins/parquet-data-format/src/test/java/org/opensearch/parquet/bridge/ParquetMergeIntegrationTests.java b/sandbox/plugins/parquet-data-format/src/test/java/org/opensearch/parquet/bridge/ParquetMergeIntegrationTests.java index cef9bd6f786f0..be07a0a9a4f31 100644 --- a/sandbox/plugins/parquet-data-format/src/test/java/org/opensearch/parquet/bridge/ParquetMergeIntegrationTests.java +++ b/sandbox/plugins/parquet-data-format/src/test/java/org/opensearch/parquet/bridge/ParquetMergeIntegrationTests.java @@ -8,6 +8,8 @@ package org.opensearch.parquet.bridge; +import com.carrotsearch.randomizedtesting.annotations.ThreadLeakScope; + import org.apache.arrow.c.ArrowArray; import org.apache.arrow.c.ArrowSchema; import org.apache.arrow.c.Data; @@ -27,6 +29,9 @@ import java.nio.file.Path; import java.util.List; +// The Tokio IO runtime worker thread (used by the Rust merge k-way merge sort) is a process-lifetime +// singleton that persists after tests complete. It polls for new async IO tasks between merges. +@ThreadLeakScope(ThreadLeakScope.Scope.NONE) public class ParquetMergeIntegrationTests extends OpenSearchTestCase { private static final String INDEX_NAME = "merge-test-index"; 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 1c0f83ceb6fa7..2652f7baa3ce7 100644 --- a/server/src/main/java/org/opensearch/index/engine/DataFormatAwareEngine.java +++ b/server/src/main/java/org/opensearch/index/engine/DataFormatAwareEngine.java @@ -261,7 +261,18 @@ public DataFormatAwareEngine(EngineConfig engineConfig) { ), registry.format(config().getIndexSettings().pluggableDataFormat()) ); - this.writerGenerationCounter = new AtomicLong(0L); + long maxGenFromCommit = 0L; + try { + List initSnapshots = committer.listCommittedSnapshots(); + if (initSnapshots.isEmpty() == false) { + for (Segment seg : initSnapshots.getLast().getSegments()) { + maxGenFromCommit = Math.max(maxGenFromCommit, seg.generation()); + } + } + } catch (IOException e) { + // Fall back to 0 on error + } + this.writerGenerationCounter = new AtomicLong(maxGenFromCommit); this.writerPool = new LockablePool<>(() -> { long gen = writerGenerationCounter.incrementAndGet(); assert gen > 0 : "writer generation must be positive but was: " + gen; @@ -461,8 +472,8 @@ private TranslogDeletionPolicy getTranslogDeletionPolicy() { @Override public Engine.IndexResult index(Engine.Index index) throws IOException { assert Objects.equals(index.uid().field(), IdFieldMapper.NAME) : index.uid().field(); - assert index.origin() == Engine.Operation.Origin.PRIMARY : "DataFormatAwareEngine only supports PRIMARY origin but got: " - + index.origin(); + assert (index.origin() == Engine.Operation.Origin.PRIMARY || index.origin() == Engine.Operation.Origin.LOCAL_TRANSLOG_RECOVERY) + : "DataFormatAwareEngine only supports PRIMARY origin but got: " + index.origin(); final boolean doThrottle = index.origin().isRecovery() == false; try (ReleasableLock ignored = readLock.acquire()) { ensureOpen(); @@ -913,11 +924,18 @@ public void flush() { @Override public boolean shouldPeriodicallyFlush() { ensureOpen(); - final long localCheckpointOfLastCommit = localCheckpointTracker.getPersistedCheckpoint(); - return translogManager.shouldPeriodicallyFlush( - localCheckpointOfLastCommit, - engineConfig.getIndexSettings().getFlushThresholdSize().getBytes() - ); + try { + Map lastCommitData = committer.getLastCommittedData(); + final long localCheckpointOfLastCommit = Long.parseLong( + lastCommitData.getOrDefault(SequenceNumbers.LOCAL_CHECKPOINT_KEY, "-1") + ); + return translogManager.shouldPeriodicallyFlush( + localCheckpointOfLastCommit, + engineConfig.getIndexSettings().getFlushThresholdSize().getBytes() + ); + } catch (IOException e) { + throw new RuntimeException(e); + } } /** Triggers a refresh to flush the indexing buffer to segments. */ diff --git a/server/src/main/java/org/opensearch/index/engine/exec/coord/CatalogSnapshot.java b/server/src/main/java/org/opensearch/index/engine/exec/coord/CatalogSnapshot.java index 309579cea1650..75de94853c279 100644 --- a/server/src/main/java/org/opensearch/index/engine/exec/coord/CatalogSnapshot.java +++ b/server/src/main/java/org/opensearch/index/engine/exec/coord/CatalogSnapshot.java @@ -64,6 +64,12 @@ public abstract class CatalogSnapshot implements Writeable, Cloneable { */ private volatile Map> filesByFormatCache; + /** + * Whether this snapshot has been committed (persisted via flush). + * Package-private — managed by {@link IndexFileDeleter}. + */ + private volatile boolean committed; + protected CatalogSnapshot(String name, long generation, long version) { this.generation = generation; this.version = version; @@ -106,6 +112,22 @@ public long getVersion() { return version; } + /** + * Marks this snapshot as committed (persisted via flush). + * Package-private — only called by {@link IndexFileDeleter} and {@link CatalogSnapshotManager}. + */ + void markCommitted() { + this.committed = true; + } + + /** + * Returns whether this snapshot was committed. + * Package-private — only called by {@link IndexFileDeleter} and {@link CatalogSnapshotManager}. + */ + boolean isCommitted() { + return committed; + } + // Package-private ref counting — only accessible within exec.coord (i.e., CatalogSnapshotManager) /** 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 index ae35ff43e3843..bbc8e7ec0bb25 100644 --- 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 @@ -359,6 +359,7 @@ public GatedConditionalCloseable acquireSnapshotForCommit() { } return new GatedConditionalCloseable<>(snapshot, () -> { try { + snapshot.markCommitted(); indexFileDeleter.onCommit(snapshot); } catch (IOException e) { throw new RuntimeException("Failed to register commit [gen=" + snapshot.getGeneration() + "]", e); diff --git a/server/src/main/java/org/opensearch/index/engine/exec/coord/IndexFileDeleter.java b/server/src/main/java/org/opensearch/index/engine/exec/coord/IndexFileDeleter.java index 38f56e850450e..802be74a0845e 100644 --- a/server/src/main/java/org/opensearch/index/engine/exec/coord/IndexFileDeleter.java +++ b/server/src/main/java/org/opensearch/index/engine/exec/coord/IndexFileDeleter.java @@ -88,6 +88,7 @@ public IndexFileDeleter( if (cs.tryIncRef() == false) { throw new IllegalStateException("Committed snapshot [gen=" + cs.getGeneration() + "] is already closed"); } + cs.markCommitted(); this.committedSnapshots.add(cs); addFileReferences(cs); } @@ -157,7 +158,7 @@ public void removeFileReferences(CatalogSnapshot snapshot) throws IOException { // Delete the commit point (segments_N) BEFORE deleting data files, // because deleteCommit may call DirectoryReader.listCommits() which // needs to read segment files that are about to be deleted. - if (commitFileManager != null) { + if (commitFileManager != null && snapshot.isCommitted()) { commitFileManager.deleteCommit(snapshot); } if (filesToDelete.isEmpty() == false) {