From e408930510d776be2c069462e3227d485c38a18b Mon Sep 17 00:00:00 2001 From: Mohit Kumar Date: Tue, 24 Mar 2026 16:57:25 +0530 Subject: [PATCH 1/3] Stab at Fixing integ tests Signed-off-by: Mohit Kumar --- .../replication/ReplicationEngine.kt | 32 ++++++++++++------- 1 file changed, 20 insertions(+), 12 deletions(-) diff --git a/src/main/kotlin/org/opensearch/replication/ReplicationEngine.kt b/src/main/kotlin/org/opensearch/replication/ReplicationEngine.kt index 94a284b21..84f0424be 100644 --- a/src/main/kotlin/org/opensearch/replication/ReplicationEngine.kt +++ b/src/main/kotlin/org/opensearch/replication/ReplicationEngine.kt @@ -11,32 +11,40 @@ package org.opensearch.replication +import org.opensearch.index.engine.DeletionStrategy +import org.opensearch.index.engine.Engine import org.opensearch.index.engine.EngineConfig +import org.opensearch.index.engine.IndexingStrategy import org.opensearch.index.engine.InternalEngine import org.opensearch.index.seqno.SequenceNumbers class ReplicationEngine(config: EngineConfig) : InternalEngine(config) { - override fun assertPrimaryIncomingSequenceNumber(origin: Operation.Origin, seqNo: Long): Boolean { - assert(origin == Operation.Origin.PRIMARY) { "Expected origin PRIMARY for replicated ops but was $origin" } - assert(seqNo != SequenceNumbers.UNASSIGNED_SEQ_NO) { "Expected valid sequence number for replicated op but was unassigned" } - return true + /** + * Replicated operations arrive with Origin.PRIMARY but must be planned using the + * non-primary path to avoid assertions in the primary planning flow (e.g. + * tryAcquireInFlightDocs) that don't apply to cross-cluster replicated ops. + */ + override fun indexingStrategyForOperation(index: Engine.Index): IndexingStrategy { + return planIndexingAsNonPrimary(index) } - override fun generateSeqNoForOperationOnPrimary(operation: Operation): Long { - check(operation.seqNo() != SequenceNumbers.UNASSIGNED_SEQ_NO) { "Expected valid sequence number for replicate op but was unassigned"} - return operation.seqNo() + override fun deletionStrategyForOperation(delete: Engine.Delete): DeletionStrategy { + return planDeletionAsNonPrimary(delete) } - override fun indexingStrategyForOperation(index: Index): IndexingStrategy { - return planIndexingAsNonPrimary(index) + override fun assertPrimaryIncomingSequenceNumber(origin: Engine.Operation.Origin, seqNo: Long): Boolean { + assert(origin == Engine.Operation.Origin.PRIMARY) { "Expected origin PRIMARY for replicated ops but was $origin" } + assert(seqNo != SequenceNumbers.UNASSIGNED_SEQ_NO) { "Expected valid sequence number for replicated op but was unassigned" } + return true } - override fun deletionStrategyForOperation(delete: Delete): DeletionStrategy { - return planDeletionAsNonPrimary(delete) + override fun generateSeqNoForOperationOnPrimary(operation: Engine.Operation): Long { + check(operation.seqNo() != SequenceNumbers.UNASSIGNED_SEQ_NO) { "Expected valid sequence number for replicate op but was unassigned"} + return operation.seqNo() } - override fun assertNonPrimaryOrigin(operation: Operation): Boolean { + override fun assertNonPrimaryOrigin(operation: Engine.Operation): Boolean { return true } } From ed33475178b74beee0204f579f748c512262cb54 Mon Sep 17 00:00:00 2001 From: Mohit Kumar Date: Tue, 24 Mar 2026 17:43:31 +0530 Subject: [PATCH 2/3] Revert "Stab at Fixing integ tests" This reverts commit 0e4b126b6adbdb4558046c3b50cf746bc5293da7. Signed-off-by: Mohit Kumar --- .../replication/ReplicationEngine.kt | 32 +++++++------------ 1 file changed, 12 insertions(+), 20 deletions(-) diff --git a/src/main/kotlin/org/opensearch/replication/ReplicationEngine.kt b/src/main/kotlin/org/opensearch/replication/ReplicationEngine.kt index 84f0424be..94a284b21 100644 --- a/src/main/kotlin/org/opensearch/replication/ReplicationEngine.kt +++ b/src/main/kotlin/org/opensearch/replication/ReplicationEngine.kt @@ -11,40 +11,32 @@ package org.opensearch.replication -import org.opensearch.index.engine.DeletionStrategy -import org.opensearch.index.engine.Engine import org.opensearch.index.engine.EngineConfig -import org.opensearch.index.engine.IndexingStrategy import org.opensearch.index.engine.InternalEngine import org.opensearch.index.seqno.SequenceNumbers class ReplicationEngine(config: EngineConfig) : InternalEngine(config) { - /** - * Replicated operations arrive with Origin.PRIMARY but must be planned using the - * non-primary path to avoid assertions in the primary planning flow (e.g. - * tryAcquireInFlightDocs) that don't apply to cross-cluster replicated ops. - */ - override fun indexingStrategyForOperation(index: Engine.Index): IndexingStrategy { - return planIndexingAsNonPrimary(index) - } - - override fun deletionStrategyForOperation(delete: Engine.Delete): DeletionStrategy { - return planDeletionAsNonPrimary(delete) - } - - override fun assertPrimaryIncomingSequenceNumber(origin: Engine.Operation.Origin, seqNo: Long): Boolean { - assert(origin == Engine.Operation.Origin.PRIMARY) { "Expected origin PRIMARY for replicated ops but was $origin" } + override fun assertPrimaryIncomingSequenceNumber(origin: Operation.Origin, seqNo: Long): Boolean { + assert(origin == Operation.Origin.PRIMARY) { "Expected origin PRIMARY for replicated ops but was $origin" } assert(seqNo != SequenceNumbers.UNASSIGNED_SEQ_NO) { "Expected valid sequence number for replicated op but was unassigned" } return true } - override fun generateSeqNoForOperationOnPrimary(operation: Engine.Operation): Long { + override fun generateSeqNoForOperationOnPrimary(operation: Operation): Long { check(operation.seqNo() != SequenceNumbers.UNASSIGNED_SEQ_NO) { "Expected valid sequence number for replicate op but was unassigned"} return operation.seqNo() } - override fun assertNonPrimaryOrigin(operation: Engine.Operation): Boolean { + override fun indexingStrategyForOperation(index: Index): IndexingStrategy { + return planIndexingAsNonPrimary(index) + } + + override fun deletionStrategyForOperation(delete: Delete): DeletionStrategy { + return planDeletionAsNonPrimary(delete) + } + + override fun assertNonPrimaryOrigin(operation: Operation): Boolean { return true } } From e9b381cd9fd0ff4ded56877133101b00d60e316f Mon Sep 17 00:00:00 2001 From: Mohit Kumar Date: Tue, 24 Mar 2026 20:31:44 +0530 Subject: [PATCH 3/3] Stab at Fixing integ tests Signed-off-by: Mohit Kumar --- .../replication/ReplicationEngine.kt | 24 +++++++++++++++++-- 1 file changed, 22 insertions(+), 2 deletions(-) diff --git a/src/main/kotlin/org/opensearch/replication/ReplicationEngine.kt b/src/main/kotlin/org/opensearch/replication/ReplicationEngine.kt index 94a284b21..f013828a2 100644 --- a/src/main/kotlin/org/opensearch/replication/ReplicationEngine.kt +++ b/src/main/kotlin/org/opensearch/replication/ReplicationEngine.kt @@ -11,7 +11,9 @@ package org.opensearch.replication +import org.opensearch.index.engine.DeletionStrategy import org.opensearch.index.engine.EngineConfig +import org.opensearch.index.engine.IndexingStrategy import org.opensearch.index.engine.InternalEngine import org.opensearch.index.seqno.SequenceNumbers @@ -28,12 +30,30 @@ class ReplicationEngine(config: EngineConfig) : InternalEngine(config) { return operation.seqNo() } + /** + * Route all index operations through the non-primary planning path. CCR replays operations with + * PRIMARY origin but they need non-primary planning to avoid version conflict checks and in-flight + * doc acquisition. We create a copy with REPLICA origin to satisfy the assertions in both the engine + * and the new OperationStrategyPlanner. + */ override fun indexingStrategyForOperation(index: Index): IndexingStrategy { - return planIndexingAsNonPrimary(index) + val replicaIndex = Index( + index.uid(), index.parsedDoc(), index.seqNo(), index.primaryTerm(), + index.version(), null, Operation.Origin.REPLICA, + index.startTime(), index.getAutoGeneratedIdTimestamp(), index.isRetry, + SequenceNumbers.UNASSIGNED_SEQ_NO, 0 + ) + return planIndexingAsNonPrimary(replicaIndex) } + // Same as above for delete operations. override fun deletionStrategyForOperation(delete: Delete): DeletionStrategy { - return planDeletionAsNonPrimary(delete) + val replicaDelete = Delete( + delete.id(), delete.uid(), delete.seqNo(), delete.primaryTerm(), + delete.version(), null, Operation.Origin.REPLICA, + delete.startTime(), SequenceNumbers.UNASSIGNED_SEQ_NO, 0 + ) + return planDeletionAsNonPrimary(replicaDelete) } override fun assertNonPrimaryOrigin(operation: Operation): Boolean {