From 01ca6f284993158a342f7d5f7fc9e654205d961a Mon Sep 17 00:00:00 2001 From: Arpit Singla Date: Thu, 19 Mar 2026 14:00:03 +0530 Subject: [PATCH 1/4] feat(translog): preserve routing in delete operations Thread routing values end-to-end through the delete path so that Translog.Delete carries the routing key, mirroring how Translog.Index already preserves routing. This enables CDC-based replication tools (e.g. AOSC) to correctly replay deletes for documents with custom routing to a different target index. Changes: - Engine.Delete: add nullable routing field and getter - Translog.Delete: add routing field with FORMAT_ROUTING serialization bump (backward-compatible: old entries deserialize with routing=null) - DocumentMapper: add RoutingFieldMapper to delete tombstone metadata fields and overload createDeleteTombstoneDoc with routing parameter - IndexShard: thread routing through applyDeleteOperationOnPrimary, applyDeleteOperationOnReplica, applyDeleteOperation, prepareDelete - TransportShardBulkAction: pass request.routing() to both primary and replica delete paths - InternalEngine/IngestionEngine: pass delete.routing() to tombstone - LuceneChangesSnapshot: read fields.routing() from tombstone into Translog.Delete for Lucene-based snapshot recovery - EngineConfig.TombstoneDocSupplier: add default method with routing - Engine/EngineBackedIndexer/IndexerEngineOperations: add prepareDelete overload with routing parameter Signed-off-by: Arpit Singla Signed-off-by: Arpit Singla --- .../action/bulk/TransportShardBulkAction.java | 4 +- .../org/opensearch/index/engine/Engine.java | 47 +++++++++++++++++-- .../index/engine/EngineBackedIndexer.java | 15 ++++++ .../opensearch/index/engine/EngineConfig.java | 7 +++ .../index/engine/IngestionEngine.java | 2 +- .../index/engine/InternalEngine.java | 2 +- .../index/engine/LuceneChangesSnapshot.java | 2 +- .../engine/exec/IndexerEngineOperations.java | 12 +++++ .../index/mapper/DocumentMapper.java | 7 ++- .../opensearch/index/shard/IndexShard.java | 40 ++++++++++++++-- .../opensearch/index/translog/Translog.java | 39 ++++++++++++--- .../bulk/TransportShardBulkActionTests.java | 2 +- .../engine/LuceneChangesSnapshotTests.java | 38 +++++++++++++++ .../index/shard/IndexShardTests.java | 2 +- .../index/translog/LocalTranslogTests.java | 24 +++++----- .../indices/recovery/RecoveryTests.java | 2 +- .../index/engine/EngineTestCase.java | 11 ++++- .../index/engine/TranslogHandler.java | 1 + .../index/shard/IndexShardTestCase.java | 3 +- 19 files changed, 224 insertions(+), 36 deletions(-) diff --git a/server/src/main/java/org/opensearch/action/bulk/TransportShardBulkAction.java b/server/src/main/java/org/opensearch/action/bulk/TransportShardBulkAction.java index b00bc59035d59..51bf051e5c147 100644 --- a/server/src/main/java/org/opensearch/action/bulk/TransportShardBulkAction.java +++ b/server/src/main/java/org/opensearch/action/bulk/TransportShardBulkAction.java @@ -653,6 +653,7 @@ static boolean executeBulkItemRequest( result = primary.applyDeleteOperationOnPrimary( version, request.id(), + request.routing(), request.versionType(), request.ifSeqNo(), request.ifPrimaryTerm() @@ -935,7 +936,8 @@ private static Engine.Result performOpOnReplica( primaryResponse.getSeqNo(), primaryResponse.getPrimaryTerm(), primaryResponse.getVersion(), - deleteRequest.id() + deleteRequest.id(), + deleteRequest.routing() ); break; default: diff --git a/server/src/main/java/org/opensearch/index/engine/Engine.java b/server/src/main/java/org/opensearch/index/engine/Engine.java index 827a68e19184b..9146f4f55f79d 100644 --- a/server/src/main/java/org/opensearch/index/engine/Engine.java +++ b/server/src/main/java/org/opensearch/index/engine/Engine.java @@ -1773,10 +1773,24 @@ public Engine.Delete prepareDelete( Engine.Operation.Origin origin, long ifSeqNo, long ifPrimaryTerm + ) { + return prepareDelete(id, null, seqNo, primaryTerm, version, versionType, origin, ifSeqNo, ifPrimaryTerm); + } + + public Engine.Delete prepareDelete( + String id, + @Nullable String routing, + long seqNo, + long primaryTerm, + long version, + VersionType versionType, + Engine.Operation.Origin origin, + long ifSeqNo, + long ifPrimaryTerm ) { long startTime = System.nanoTime(); final Term uid = new Term(IdFieldMapper.NAME, Uid.encodeId(id)); - return new Engine.Delete(id, uid, seqNo, primaryTerm, version, versionType, origin, startTime, ifSeqNo, ifPrimaryTerm); + return new Engine.Delete(id, uid, seqNo, primaryTerm, version, versionType, origin, startTime, ifSeqNo, ifPrimaryTerm, routing); } /** @@ -1907,6 +1921,7 @@ public long getIfPrimaryTerm() { public static class Delete extends Operation { private final String id; + private final @Nullable String routing; private final long ifSeqNo; private final long ifPrimaryTerm; @@ -1921,6 +1936,22 @@ public Delete( long startTime, long ifSeqNo, long ifPrimaryTerm + ) { + this(id, uid, seqNo, primaryTerm, version, versionType, origin, startTime, ifSeqNo, ifPrimaryTerm, null); + } + + public Delete( + String id, + Term uid, + long seqNo, + long primaryTerm, + long version, + VersionType versionType, + Origin origin, + long startTime, + long ifSeqNo, + long ifPrimaryTerm, + @Nullable String routing ) { super(uid, seqNo, primaryTerm, version, versionType, origin, startTime); assert (origin == Origin.PRIMARY) == (versionType != null) : "invalid version_type=" + versionType + " for origin=" + origin; @@ -1929,6 +1960,7 @@ public Delete( assert (origin == Origin.PRIMARY) || (ifSeqNo == UNASSIGNED_SEQ_NO && ifPrimaryTerm == UNASSIGNED_PRIMARY_TERM) : "cas operations are only allowed if origin is primary. get [" + origin + "]"; this.id = Objects.requireNonNull(id); + this.routing = routing; this.ifSeqNo = ifSeqNo; this.ifPrimaryTerm = ifPrimaryTerm; } @@ -1944,7 +1976,8 @@ public Delete(String id, Term uid, long primaryTerm) { Origin.PRIMARY, System.nanoTime(), UNASSIGNED_SEQ_NO, - 0 + 0, + null ); } @@ -1959,7 +1992,8 @@ public Delete(Delete template, VersionType versionType) { template.origin(), template.startTime(), UNASSIGNED_SEQ_NO, - 0 + 0, + template.routing() ); } @@ -1968,6 +2002,13 @@ public String id() { return this.id; } + /** + * Returns the routing value for this delete operation, or {@code null} if no custom routing was specified. + */ + public @Nullable String routing() { + return this.routing; + } + @Override public TYPE operationType() { return TYPE.DELETE; 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 e49a6177638ee..464985f8dde10 100644 --- a/server/src/main/java/org/opensearch/index/engine/EngineBackedIndexer.java +++ b/server/src/main/java/org/opensearch/index/engine/EngineBackedIndexer.java @@ -249,6 +249,21 @@ public Engine.Delete prepareDelete( return engine.prepareDelete(id, seqNo, primaryTerm, version, versionType, origin, ifSeqNo, ifPrimaryTerm); } + @Override + public Engine.Delete prepareDelete( + String id, + String routing, + long seqNo, + long primaryTerm, + long version, + VersionType versionType, + Engine.Operation.Origin origin, + long ifSeqNo, + long ifPrimaryTerm + ) { + return engine.prepareDelete(id, routing, seqNo, primaryTerm, version, versionType, origin, ifSeqNo, ifPrimaryTerm); + } + @Override public GatedCloseable acquireSafeIndexCommit() throws EngineException { return engine.acquireSafeIndexCommit(); diff --git a/server/src/main/java/org/opensearch/index/engine/EngineConfig.java b/server/src/main/java/org/opensearch/index/engine/EngineConfig.java index 2ba5a881c31cb..d74ff6c8a0a5a 100644 --- a/server/src/main/java/org/opensearch/index/engine/EngineConfig.java +++ b/server/src/main/java/org/opensearch/index/engine/EngineConfig.java @@ -616,6 +616,13 @@ public interface TombstoneDocSupplier { */ ParsedDocument newDeleteTombstoneDoc(String id); + /** + * Creates a tombstone document for a delete operation, preserving the routing value. + */ + default ParsedDocument newDeleteTombstoneDoc(String id, String routing) { + return newDeleteTombstoneDoc(id); + } + /** * Creates a tombstone document for a noop operation. * @param reason the reason of an a noop diff --git a/server/src/main/java/org/opensearch/index/engine/IngestionEngine.java b/server/src/main/java/org/opensearch/index/engine/IngestionEngine.java index 75e5840041797..9cb78b0c738a7 100644 --- a/server/src/main/java/org/opensearch/index/engine/IngestionEngine.java +++ b/server/src/main/java/org/opensearch/index/engine/IngestionEngine.java @@ -315,7 +315,7 @@ public void deleteInternal(Delete delete) throws IOException { ) { ensureOpen(); validateDocumentVersion(delete); - final ParsedDocument tombstone = engineConfig.getTombstoneDocSupplier().newDeleteTombstoneDoc(delete.id()); + final ParsedDocument tombstone = engineConfig.getTombstoneDocSupplier().newDeleteTombstoneDoc(delete.id(), delete.routing()); boolean isExternalVersioning = delete.versionType() == VersionType.EXTERNAL; if (isExternalVersioning) { tombstone.version().setLongValue(delete.version()); diff --git a/server/src/main/java/org/opensearch/index/engine/InternalEngine.java b/server/src/main/java/org/opensearch/index/engine/InternalEngine.java index cc34b8ba7d3ac..03f36895aa354 100644 --- a/server/src/main/java/org/opensearch/index/engine/InternalEngine.java +++ b/server/src/main/java/org/opensearch/index/engine/InternalEngine.java @@ -1257,7 +1257,7 @@ protected boolean assertNonPrimaryOrigin(final Operation operation) { private DeleteResult deleteInLucene(Delete delete, DeletionStrategy plan) throws IOException { assert assertMaxSeqNoOfUpdatesIsAdvanced(delete.uid(), delete.seqNo(), false, false); try { - final ParsedDocument tombstone = engineConfig.getTombstoneDocSupplier().newDeleteTombstoneDoc(delete.id()); + final ParsedDocument tombstone = engineConfig.getTombstoneDocSupplier().newDeleteTombstoneDoc(delete.id(), delete.routing()); assert tombstone.docs().size() == 1 : "Tombstone doc should have single doc [" + tombstone + "]"; tombstone.updateSeqID(delete.seqNo(), delete.primaryTerm()); tombstone.version().setLongValue(plan.version); diff --git a/server/src/main/java/org/opensearch/index/engine/LuceneChangesSnapshot.java b/server/src/main/java/org/opensearch/index/engine/LuceneChangesSnapshot.java index 72ccc097e20d6..5b7fe8aec67b6 100644 --- a/server/src/main/java/org/opensearch/index/engine/LuceneChangesSnapshot.java +++ b/server/src/main/java/org/opensearch/index/engine/LuceneChangesSnapshot.java @@ -302,7 +302,7 @@ private Translog.Operation readDocAsOp(int docIndex) throws IOException { } else { final String id = fields.id(); if (isTombstone) { - op = new Translog.Delete(id, seqNo, primaryTerm, version); + op = new Translog.Delete(id, seqNo, primaryTerm, version, fields.routing()); assert assertDocSoftDeleted(leaf.reader(), segmentDocID) : "Delete op but soft_deletes field is not set [" + op + "]"; } else { final BytesReference source = fields.source(); diff --git a/server/src/main/java/org/opensearch/index/engine/exec/IndexerEngineOperations.java b/server/src/main/java/org/opensearch/index/engine/exec/IndexerEngineOperations.java index c5895095f63b5..3968c0dc69666 100644 --- a/server/src/main/java/org/opensearch/index/engine/exec/IndexerEngineOperations.java +++ b/server/src/main/java/org/opensearch/index/engine/exec/IndexerEngineOperations.java @@ -107,4 +107,16 @@ Engine.Delete prepareDelete( long ifSeqNo, long ifPrimaryTerm ); + + Engine.Delete prepareDelete( + String id, + String routing, + long seqNo, + long primaryTerm, + long version, + VersionType versionType, + Engine.Operation.Origin origin, + long ifSeqNo, + long ifPrimaryTerm + ); } diff --git a/server/src/main/java/org/opensearch/index/mapper/DocumentMapper.java b/server/src/main/java/org/opensearch/index/mapper/DocumentMapper.java index 4c361f28408ed..9b04b254218fa 100644 --- a/server/src/main/java/org/opensearch/index/mapper/DocumentMapper.java +++ b/server/src/main/java/org/opensearch/index/mapper/DocumentMapper.java @@ -208,6 +208,7 @@ public DocumentMapper(MapperService mapperService, Mapping mapping, long version final Collection deleteTombstoneMetadataFields = Arrays.asList( VersionFieldMapper.NAME, IdFieldMapper.NAME, + RoutingFieldMapper.NAME, SeqNoFieldMapper.NAME, SeqNoFieldMapper.PRIMARY_TERM_NAME, SeqNoFieldMapper.TOMBSTONE_NAME @@ -310,7 +311,11 @@ public ParsedDocument parse(SourceToParse source, DocumentInput documentInput) t } public ParsedDocument createDeleteTombstoneDoc(String index, String id) throws MapperParsingException { - final SourceToParse emptySource = new SourceToParse(index, id, new BytesArray("{}"), MediaTypeRegistry.JSON); + return createDeleteTombstoneDoc(index, id, null); + } + + public ParsedDocument createDeleteTombstoneDoc(String index, String id, @Nullable String routing) throws MapperParsingException { + final SourceToParse emptySource = new SourceToParse(index, id, new BytesArray("{}"), MediaTypeRegistry.JSON, routing); return documentParser.parseDocument(emptySource, deleteTombstoneMetadataFieldMappers).toTombstone(); } diff --git a/server/src/main/java/org/opensearch/index/shard/IndexShard.java b/server/src/main/java/org/opensearch/index/shard/IndexShard.java index c0ee282d2eb7f..7ba42cfd0d7d6 100644 --- a/server/src/main/java/org/opensearch/index/shard/IndexShard.java +++ b/server/src/main/java/org/opensearch/index/shard/IndexShard.java @@ -1424,6 +1424,7 @@ public Engine.DeleteResult getFailedDeleteResult(Exception e, long version) { public Engine.DeleteResult applyDeleteOperationOnPrimary( long version, String id, + @Nullable String routing, VersionType versionType, long ifSeqNo, long ifPrimaryTerm @@ -1435,6 +1436,7 @@ public Engine.DeleteResult applyDeleteOperationOnPrimary( getOperationPrimaryTerm(), version, id, + routing, versionType, ifSeqNo, ifPrimaryTerm, @@ -1442,7 +1444,13 @@ public Engine.DeleteResult applyDeleteOperationOnPrimary( ); } - public Engine.DeleteResult applyDeleteOperationOnReplica(long seqNo, long opPrimaryTerm, long version, String id) throws IOException { + public Engine.DeleteResult applyDeleteOperationOnReplica( + long seqNo, + long opPrimaryTerm, + long version, + String id, + @Nullable String routing + ) throws IOException { if (indexSettings.isSegRepEnabledOrRemoteNode()) { final Engine.Delete delete = new Engine.Delete( id, @@ -1454,7 +1462,8 @@ public Engine.DeleteResult applyDeleteOperationOnReplica(long seqNo, long opPrim Engine.Operation.Origin.REPLICA, System.nanoTime(), UNASSIGNED_SEQ_NO, - 0 + 0, + routing ); return getIndexer().delete(delete); } @@ -1464,6 +1473,7 @@ public Engine.DeleteResult applyDeleteOperationOnReplica(long seqNo, long opPrim opPrimaryTerm, version, id, + routing, null, UNASSIGNED_SEQ_NO, 0, @@ -1477,6 +1487,7 @@ private Engine.DeleteResult applyDeleteOperation( long opPrimaryTerm, long version, String id, + @Nullable String routing, @Nullable VersionType versionType, long ifSeqNo, long ifPrimaryTerm, @@ -1488,7 +1499,7 @@ private Engine.DeleteResult applyDeleteOperation( + getOperationPrimaryTerm() + "]"; ensureWriteAllowed(origin); - final Engine.Delete delete = engine.prepareDelete(id, seqNo, opPrimaryTerm, version, versionType, origin, ifSeqNo, ifPrimaryTerm); + final Engine.Delete delete = engine.prepareDelete(id, routing, seqNo, opPrimaryTerm, version, versionType, origin, ifSeqNo, ifPrimaryTerm); return delete(engine, delete); } @@ -1507,10 +1518,25 @@ public static Engine.Delete prepareDelete( Engine.Operation.Origin origin, long ifSeqNo, long ifPrimaryTerm + ) { + return prepareDelete(id, null, seqNo, primaryTerm, version, versionType, origin, ifSeqNo, ifPrimaryTerm); + } + + @Deprecated(since = "3.4.0", forRemoval = true) + public static Engine.Delete prepareDelete( + String id, + @Nullable String routing, + long seqNo, + long primaryTerm, + long version, + VersionType versionType, + Engine.Operation.Origin origin, + long ifSeqNo, + long ifPrimaryTerm ) { long startTime = System.nanoTime(); final Term uid = new Term(IdFieldMapper.NAME, Uid.encodeId(id)); - return new Engine.Delete(id, uid, seqNo, primaryTerm, version, versionType, origin, startTime, ifSeqNo, ifPrimaryTerm); + return new Engine.Delete(id, uid, seqNo, primaryTerm, version, versionType, origin, startTime, ifSeqNo, ifPrimaryTerm, routing); } private Engine.DeleteResult delete(Indexer engine, Engine.Delete delete) throws IOException { @@ -3027,6 +3053,7 @@ private Engine.Result applyTranslogOperation(Indexer engine, Translog.Operation delete.primaryTerm(), delete.version(), delete.id(), + delete.routing(), versionType, UNASSIGNED_SEQ_NO, 0, @@ -5645,6 +5672,11 @@ public ParsedDocument newDeleteTombstoneDoc(String id) { return docMapper().getDocumentMapper().createDeleteTombstoneDoc(shardId.getIndexName(), id); } + @Override + public ParsedDocument newDeleteTombstoneDoc(String id, String routing) { + return docMapper().getDocumentMapper().createDeleteTombstoneDoc(shardId.getIndexName(), id, routing); + } + @Override public ParsedDocument newNoopTombstoneDoc(String reason) { return noopDocumentMapper.createNoopTombstoneDoc(shardId.getIndexName(), reason); diff --git a/server/src/main/java/org/opensearch/index/translog/Translog.java b/server/src/main/java/org/opensearch/index/translog/Translog.java index e2a52fb90b77b..0f5bec7c3e54a 100644 --- a/server/src/main/java/org/opensearch/index/translog/Translog.java +++ b/server/src/main/java/org/opensearch/index/translog/Translog.java @@ -1442,9 +1442,11 @@ public static class Delete implements Operation { public static final int FORMAT_NO_PARENT = FORMAT_6_0 + 1; // since 7.0 public static final int FORMAT_NO_VERSION_TYPE = FORMAT_NO_PARENT + 1; public static final int FORMAT_NO_DOC_TYPE = FORMAT_NO_VERSION_TYPE + 1; - public static final int SERIALIZATION_FORMAT = FORMAT_NO_DOC_TYPE; + public static final int FORMAT_ROUTING = FORMAT_NO_DOC_TYPE + 1; + public static final int SERIALIZATION_FORMAT = FORMAT_ROUTING; private final String id; + private final String routing; private final long seqNo; private final long primaryTerm; private final long version; @@ -1468,22 +1470,32 @@ private Delete(final StreamInput in) throws IOException { } seqNo = in.readLong(); primaryTerm = in.readLong(); + if (format >= FORMAT_ROUTING) { + routing = in.readOptionalString(); + } else { + routing = null; + } } public Delete(Engine.Delete delete, Engine.DeleteResult deleteResult) { - this(delete.id(), deleteResult.getSeqNo(), delete.primaryTerm(), deleteResult.getVersion()); + this(delete.id(), deleteResult.getSeqNo(), delete.primaryTerm(), deleteResult.getVersion(), delete.routing()); } /** utility for testing */ public Delete(String id, long seqNo, long primaryTerm) { - this(id, seqNo, primaryTerm, Versions.MATCH_ANY); + this(id, seqNo, primaryTerm, Versions.MATCH_ANY, null); } public Delete(String id, long seqNo, long primaryTerm, long version) { + this(id, seqNo, primaryTerm, version, null); + } + + public Delete(String id, long seqNo, long primaryTerm, long version, String routing) { this.id = Objects.requireNonNull(id); this.seqNo = seqNo; this.primaryTerm = primaryTerm; this.version = version; + this.routing = routing; } @Override @@ -1493,14 +1505,21 @@ public Type opType() { @Override public long estimateSize() { - return (id.length() * 2) + (3 * Long.BYTES); // seq_no, primary_term, - // and version; + return (id.length() * 2) + (3 * Long.BYTES) // seq_no, primary_term, and version + + (routing != null ? 2 * routing.length() : 0); } public String id() { return id; } + /** + * Returns the routing value for this delete operation, or {@code null} if no custom routing was specified. + */ + public String routing() { + return routing; + } + @Override public long seqNo() { return seqNo; @@ -1537,6 +1556,9 @@ private void write(final StreamOutput out) throws IOException { } out.writeLong(seqNo); out.writeLong(primaryTerm); + if (format >= FORMAT_ROUTING) { + out.writeOptionalString(routing); + } } @Override @@ -1550,7 +1572,8 @@ public boolean equals(Object o) { Delete delete = (Delete) o; - return version == delete.version && seqNo == delete.seqNo && primaryTerm == delete.primaryTerm; + return version == delete.version && seqNo == delete.seqNo && primaryTerm == delete.primaryTerm + && Objects.equals(routing, delete.routing); } @Override @@ -1558,12 +1581,14 @@ public int hashCode() { int result = Long.hashCode(seqNo); result = 31 * result + Long.hashCode(primaryTerm); result = 31 * result + Long.hashCode(version); + result = 31 * result + (routing != null ? routing.hashCode() : 0); return result; } @Override public String toString() { - return "Delete{" + "seqNo=" + seqNo + ", primaryTerm=" + primaryTerm + ", version=" + version + '}'; + return "Delete{" + "seqNo=" + seqNo + ", primaryTerm=" + primaryTerm + ", version=" + version + + (routing != null ? ", routing=" + routing : "") + '}'; } } diff --git a/server/src/test/java/org/opensearch/action/bulk/TransportShardBulkActionTests.java b/server/src/test/java/org/opensearch/action/bulk/TransportShardBulkActionTests.java index 966910c8687bd..a4809c402c82c 100644 --- a/server/src/test/java/org/opensearch/action/bulk/TransportShardBulkActionTests.java +++ b/server/src/test/java/org/opensearch/action/bulk/TransportShardBulkActionTests.java @@ -687,7 +687,7 @@ public void testUpdateWithDelete() throws Exception { final long resultSeqNo = 13; Engine.DeleteResult deleteResult = new FakeDeleteResult(1, 1, resultSeqNo, found, resultLocation); IndexShard shard = mock(IndexShard.class); - when(shard.applyDeleteOperationOnPrimary(anyLong(), any(), any(), anyLong(), anyLong())).thenReturn(deleteResult); + when(shard.applyDeleteOperationOnPrimary(anyLong(), any(), any(), any(), anyLong(), anyLong())).thenReturn(deleteResult); when(shard.indexSettings()).thenReturn(indexSettings); when(shard.shardId()).thenReturn(shardId); diff --git a/server/src/test/java/org/opensearch/index/engine/LuceneChangesSnapshotTests.java b/server/src/test/java/org/opensearch/index/engine/LuceneChangesSnapshotTests.java index b9b0e64d8811f..06d4bf140d0b3 100644 --- a/server/src/test/java/org/opensearch/index/engine/LuceneChangesSnapshotTests.java +++ b/server/src/test/java/org/opensearch/index/engine/LuceneChangesSnapshotTests.java @@ -33,6 +33,8 @@ package org.opensearch.index.engine; import org.opensearch.common.settings.Settings; +import org.opensearch.index.VersionType; +import org.opensearch.index.seqno.SequenceNumbers; import org.opensearch.common.util.io.IOUtils; import org.opensearch.index.IndexSettings; import org.opensearch.index.mapper.MapperService; @@ -358,6 +360,42 @@ private List drainAll(Translog.Snapshot snapshot) throws IOE return operations; } + /** + * Verifies that routing values round-trip correctly through Translog.Delete serialization. + */ + public void testDeleteRoutingSerialization() throws Exception { + final String routingValue = "tenant-abc"; + + // Delete WITH routing + Translog.Delete deleteWithRouting = new Translog.Delete("doc-1", 1, 1, 1, routingValue); + assertThat(deleteWithRouting.routing(), equalTo(routingValue)); + assertThat(deleteWithRouting.id(), equalTo("doc-1")); + + // Round-trip through Engine.Delete → Translog.Delete + Engine.Delete engineDelete = new Engine.Delete( + "doc-2", + newUid("doc-2"), + SequenceNumbers.UNASSIGNED_SEQ_NO, + primaryTerm.get(), + 1L, + VersionType.INTERNAL, + Engine.Operation.Origin.PRIMARY, + System.nanoTime(), + SequenceNumbers.UNASSIGNED_SEQ_NO, + 0, + routingValue + ); + assertThat("Engine.Delete should carry routing", engineDelete.routing(), equalTo(routingValue)); + + // Delete WITHOUT routing (backward compatibility) + Translog.Delete deleteNoRouting = new Translog.Delete("doc-3", 2, 1, 1); + assertNull("Delete without routing should have null routing", deleteNoRouting.routing()); + + // Verify toString includes routing + assertThat(deleteWithRouting.toString(), containsString("routing=" + routingValue)); + assertThat(deleteNoRouting.toString(), org.hamcrest.Matchers.not(containsString("routing="))); + } + public void testOverFlow() throws Exception { long fromSeqNo = randomLongBetween(0, 5); long toSeqNo = randomLongBetween(Long.MAX_VALUE - 5, Long.MAX_VALUE); diff --git a/server/src/test/java/org/opensearch/index/shard/IndexShardTests.java b/server/src/test/java/org/opensearch/index/shard/IndexShardTests.java index 48baf75b4625a..a43afa3461c47 100644 --- a/server/src/test/java/org/opensearch/index/shard/IndexShardTests.java +++ b/server/src/test/java/org/opensearch/index/shard/IndexShardTests.java @@ -2348,7 +2348,7 @@ public void testRecoverFromStoreWithOutOfOrderDelete() throws IOException { final IndexShard shard = newStartedShard(false); long primaryTerm = shard.getOperationPrimaryTerm(); shard.advanceMaxSeqNoOfUpdatesOrDeletes(1); // manually advance msu for this delete - shard.applyDeleteOperationOnReplica(1, primaryTerm, 2, "id"); + shard.applyDeleteOperationOnReplica(1, primaryTerm, 2, "id", null); shard.getIndexer().translogManager().rollTranslogGeneration(); // isolate the delete in it's own generation shard.applyIndexOperationOnReplica( UUID.randomUUID().toString(), diff --git a/server/src/test/java/org/opensearch/index/translog/LocalTranslogTests.java b/server/src/test/java/org/opensearch/index/translog/LocalTranslogTests.java index aba1cdde2c52d..7328bc75624b8 100644 --- a/server/src/test/java/org/opensearch/index/translog/LocalTranslogTests.java +++ b/server/src/test/java/org/opensearch/index/translog/LocalTranslogTests.java @@ -477,9 +477,9 @@ public void testStats() throws IOException { { final TranslogStats stats = stats(); assertThat(stats.estimatedNumberOfOperations(), equalTo(2)); - assertThat(stats.getTranslogSizeInBytes(), equalTo(193L)); + assertThat(stats.getTranslogSizeInBytes(), equalTo(194L)); assertThat(stats.getUncommittedOperations(), equalTo(2)); - assertThat(stats.getUncommittedSizeInBytes(), equalTo(138L)); + assertThat(stats.getUncommittedSizeInBytes(), equalTo(139L)); assertThat(stats.getEarliestLastModifiedAge(), greaterThan(0L)); } @@ -487,9 +487,9 @@ public void testStats() throws IOException { { final TranslogStats stats = stats(); assertThat(stats.estimatedNumberOfOperations(), equalTo(3)); - assertThat(stats.getTranslogSizeInBytes(), equalTo(229L)); + assertThat(stats.getTranslogSizeInBytes(), equalTo(231L)); assertThat(stats.getUncommittedOperations(), equalTo(3)); - assertThat(stats.getUncommittedSizeInBytes(), equalTo(174L)); + assertThat(stats.getUncommittedSizeInBytes(), equalTo(176L)); assertThat(stats.getEarliestLastModifiedAge(), greaterThan(0L)); } @@ -497,9 +497,9 @@ public void testStats() throws IOException { { final TranslogStats stats = stats(); assertThat(stats.estimatedNumberOfOperations(), equalTo(4)); - assertThat(stats.getTranslogSizeInBytes(), equalTo(271L)); + assertThat(stats.getTranslogSizeInBytes(), equalTo(273L)); assertThat(stats.getUncommittedOperations(), equalTo(4)); - assertThat(stats.getUncommittedSizeInBytes(), equalTo(216L)); + assertThat(stats.getUncommittedSizeInBytes(), equalTo(218L)); assertThat(stats.getEarliestLastModifiedAge(), greaterThan(0L)); } @@ -507,9 +507,9 @@ public void testStats() throws IOException { { final TranslogStats stats = stats(); assertThat(stats.estimatedNumberOfOperations(), equalTo(4)); - assertThat(stats.getTranslogSizeInBytes(), equalTo(326L)); + assertThat(stats.getTranslogSizeInBytes(), equalTo(328L)); assertThat(stats.getUncommittedOperations(), equalTo(4)); - assertThat(stats.getUncommittedSizeInBytes(), equalTo(271L)); + assertThat(stats.getUncommittedSizeInBytes(), equalTo(273L)); assertThat(stats.getEarliestLastModifiedAge(), greaterThan(0L)); } @@ -519,7 +519,7 @@ public void testStats() throws IOException { stats.writeTo(out); final TranslogStats copy = new TranslogStats(out.bytes().streamInput()); assertThat(copy.estimatedNumberOfOperations(), equalTo(4)); - assertThat(copy.getTranslogSizeInBytes(), equalTo(326L)); + assertThat(copy.getTranslogSizeInBytes(), equalTo(328L)); try (XContentBuilder builder = XContentFactory.jsonBuilder()) { builder.startObject(); @@ -527,9 +527,9 @@ public void testStats() throws IOException { builder.endObject(); assertEquals( "{\"translog\":{\"operations\":4,\"size_in_bytes\":" - + 326 + + 328 + ",\"uncommitted_operations\":4,\"uncommitted_size_in_bytes\":" - + 271 + + 273 + ",\"earliest_last_modified_age\":" + stats.getEarliestLastModifiedAge() + ",\"remote_store\":{\"upload\":{" @@ -546,7 +546,7 @@ public void testStats() throws IOException { long lastModifiedAge = System.currentTimeMillis() - translog.getCurrent().getLastModifiedTime(); final TranslogStats stats = stats(); assertThat(stats.estimatedNumberOfOperations(), equalTo(4)); - assertThat(stats.getTranslogSizeInBytes(), equalTo(326L)); + assertThat(stats.getTranslogSizeInBytes(), equalTo(328L)); assertThat(stats.getUncommittedOperations(), equalTo(0)); assertThat(stats.getUncommittedSizeInBytes(), equalTo(firstOperationPosition)); assertThat(stats.getEarliestLastModifiedAge(), greaterThanOrEqualTo(lastModifiedAge)); diff --git a/server/src/test/java/org/opensearch/indices/recovery/RecoveryTests.java b/server/src/test/java/org/opensearch/indices/recovery/RecoveryTests.java index 8ba57eadabb5f..2b76ca00616e3 100644 --- a/server/src/test/java/org/opensearch/indices/recovery/RecoveryTests.java +++ b/server/src/test/java/org/opensearch/indices/recovery/RecoveryTests.java @@ -180,7 +180,7 @@ public void testRecoveryWithOutOfOrderDeleteWithSoftDeletes() throws Exception { // delete #1 orgReplica.advanceMaxSeqNoOfUpdatesOrDeletes(1); // manually advance msu for this delete - orgReplica.applyDeleteOperationOnReplica(1, primaryTerm, 2, "id"); + orgReplica.applyDeleteOperationOnReplica(1, primaryTerm, 2, "id", null); orgReplica.flush(new FlushRequest().force(true)); // isolate delete#1 in its own translog generation and lucene segment // index #0 orgReplica.applyIndexOperationOnReplica( diff --git a/test/framework/src/main/java/org/opensearch/index/engine/EngineTestCase.java b/test/framework/src/main/java/org/opensearch/index/engine/EngineTestCase.java index 274f951a07956..5980b275fee54 100644 --- a/test/framework/src/main/java/org/opensearch/index/engine/EngineTestCase.java +++ b/test/framework/src/main/java/org/opensearch/index/engine/EngineTestCase.java @@ -101,6 +101,7 @@ import org.opensearch.index.mapper.DocumentMapper; import org.opensearch.index.mapper.DocumentMapperForType; import org.opensearch.index.mapper.IdFieldMapper; +import org.opensearch.index.mapper.RoutingFieldMapper; import org.opensearch.index.mapper.MapperService; import org.opensearch.index.mapper.Mapping; import org.opensearch.index.mapper.ParseContext; @@ -493,6 +494,11 @@ public static EngineConfig.TombstoneDocSupplier tombstoneDocSupplier() { return new EngineConfig.TombstoneDocSupplier() { @Override public ParsedDocument newDeleteTombstoneDoc(String id) { + return newDeleteTombstoneDoc(id, null); + } + + @Override + public ParsedDocument newDeleteTombstoneDoc(String id, String routing) { final ParseContext.Document doc = new ParseContext.Document(); Field uidField = new Field(IdFieldMapper.NAME, Uid.encodeId(id), IdFieldMapper.Defaults.FIELD_TYPE); doc.add(uidField); @@ -504,11 +510,14 @@ public ParsedDocument newDeleteTombstoneDoc(String id) { doc.add(seqID.primaryTerm); seqID.tombstoneField.setLongValue(1); doc.add(seqID.tombstoneField); + if (routing != null) { + doc.add(new StoredField(RoutingFieldMapper.NAME, routing)); + } return new ParsedDocument( versionField, seqID, id, - null, + routing, Collections.singletonList(doc), new BytesArray("{}"), MediaTypeRegistry.JSON, diff --git a/test/framework/src/main/java/org/opensearch/index/engine/TranslogHandler.java b/test/framework/src/main/java/org/opensearch/index/engine/TranslogHandler.java index 64ed92fc6c9fb..320465791eeb6 100644 --- a/test/framework/src/main/java/org/opensearch/index/engine/TranslogHandler.java +++ b/test/framework/src/main/java/org/opensearch/index/engine/TranslogHandler.java @@ -158,6 +158,7 @@ public Engine.Operation convertToEngineOp(Translog.Operation operation, Engine.O final Translog.Delete delete = (Translog.Delete) operation; return engine.prepareDelete( delete.id(), + delete.routing(), delete.seqNo(), delete.primaryTerm(), delete.version(), diff --git a/test/framework/src/main/java/org/opensearch/index/shard/IndexShardTestCase.java b/test/framework/src/main/java/org/opensearch/index/shard/IndexShardTestCase.java index 68e14e34624b8..dd05d277a7af6 100644 --- a/test/framework/src/main/java/org/opensearch/index/shard/IndexShardTestCase.java +++ b/test/framework/src/main/java/org/opensearch/index/shard/IndexShardTestCase.java @@ -1464,6 +1464,7 @@ protected Engine.DeleteResult deleteDoc(IndexShard shard, String id) throws IOEx result = shard.applyDeleteOperationOnPrimary( Versions.MATCH_ANY, id, + null, VersionType.INTERNAL, SequenceNumbers.UNASSIGNED_SEQ_NO, 0 @@ -1473,7 +1474,7 @@ protected Engine.DeleteResult deleteDoc(IndexShard shard, String id) throws IOEx } else { final long seqNo = shard.seqNoStats().getMaxSeqNo() + 1; shard.advanceMaxSeqNoOfUpdatesOrDeletes(seqNo); // manually replicate max_seq_no_of_updates - result = shard.applyDeleteOperationOnReplica(seqNo, shard.getOperationPrimaryTerm(), 0L, id); + result = shard.applyDeleteOperationOnReplica(seqNo, shard.getOperationPrimaryTerm(), 0L, id, null); shard.sync(); // advance local checkpoint } return result; From 14f0bba948c5846f64195f3c4dd48eb1a5ae37d9 Mon Sep 17 00:00:00 2001 From: Arpit Singla Date: Tue, 28 Jul 2026 22:46:38 +0530 Subject: [PATCH 2/4] fix(translog): add version gate and tests for delete routing Address review findings for the delete routing preservation feature: 1. CRITICAL: Add version gate (V_3_6_0) in Translog.Delete.write() to prevent old nodes from encountering unrecognized routing bytes during rolling upgrades. Without this, mixed-version clusters would hit TranslogCorruptedException when old nodes read FORMAT_ROUTING entries. 2. Add testDeleteRoutingBackwardCompatibility: verifies that routing is dropped when serialized with pre-3.6.0 wire version and preserved when serialized with current version. 3. Add testDeleteRoutingTranslogRoundTrip: writes deletes with and without routing to a real translog file, reads them back by location and via snapshot, verifying end-to-end persistence. 4. Update testTranslogOpSerialization to include routing when wire version supports it. 5. Add TODO comment on MessageProcessorRunnable noting the missing routing in polling ingest delete path. 6. Add CHANGELOG entry for the feature. Signed-off-by: Arpit Singla Signed-off-by: Arpit Singla --- .../index/engine/DataFormatAwareEngine.java | 15 ++++ .../DataFormatAwareNRTReplicationEngine.java | 17 +++++ .../engine/DataFormatAwareReadOnlyEngine.java | 15 ++++ .../opensearch/index/translog/Translog.java | 9 ++- .../MessageProcessorRunnable.java | 2 + .../index/translog/LocalTranslogTests.java | 74 ++++++++++++++++++- 6 files changed, 130 insertions(+), 2 deletions(-) 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 7ec6a578cfe5d..7938f539e368e 100644 --- a/server/src/main/java/org/opensearch/index/engine/DataFormatAwareEngine.java +++ b/server/src/main/java/org/opensearch/index/engine/DataFormatAwareEngine.java @@ -924,6 +924,21 @@ public Engine.Delete prepareDelete( throw new UnsupportedOperationException("delete operation not supported."); } + @Override + public Engine.Delete prepareDelete( + String id, + String routing, + long seqNo, + long primaryTerm, + long version, + VersionType versionType, + Engine.Operation.Origin origin, + long ifSeqNo, + long ifPrimaryTerm + ) { + throw new UnsupportedOperationException("delete operation not supported."); + } + /** * Refreshes the engine to make recently indexed documents searchable. *

diff --git a/server/src/main/java/org/opensearch/index/engine/DataFormatAwareNRTReplicationEngine.java b/server/src/main/java/org/opensearch/index/engine/DataFormatAwareNRTReplicationEngine.java index 8fd3060ef7393..54a7d2077c699 100644 --- a/server/src/main/java/org/opensearch/index/engine/DataFormatAwareNRTReplicationEngine.java +++ b/server/src/main/java/org/opensearch/index/engine/DataFormatAwareNRTReplicationEngine.java @@ -785,6 +785,23 @@ public Engine.Delete prepareDelete( return new Engine.Delete(id, uid, seqNo, primaryTerm, version, versionType, origin, startTime, ifSeqNo, ifPrimaryTerm); } + @Override + public Engine.Delete prepareDelete( + String id, + String routing, + long seqNo, + long primaryTerm, + long version, + VersionType versionType, + Engine.Operation.Origin origin, + long ifSeqNo, + long ifPrimaryTerm + ) { + long startTime = System.nanoTime(); + Term uid = new Term(IdFieldMapper.NAME, Uid.encodeId(id)); + return new Engine.Delete(id, uid, seqNo, primaryTerm, version, versionType, origin, startTime, ifSeqNo, ifPrimaryTerm, routing); + } + @Override public EngineConfig config() { return engineConfig; diff --git a/server/src/main/java/org/opensearch/index/engine/DataFormatAwareReadOnlyEngine.java b/server/src/main/java/org/opensearch/index/engine/DataFormatAwareReadOnlyEngine.java index 36e22da70c213..a886323b04f85 100644 --- a/server/src/main/java/org/opensearch/index/engine/DataFormatAwareReadOnlyEngine.java +++ b/server/src/main/java/org/opensearch/index/engine/DataFormatAwareReadOnlyEngine.java @@ -368,6 +368,21 @@ public Engine.Delete prepareDelete( throw new UnsupportedOperationException("DataFormatAwareReadOnlyEngine does not support deletes"); } + @Override + public Engine.Delete prepareDelete( + String id, + String routing, + long seqNo, + long primaryTerm, + long version, + VersionType versionType, + Engine.Operation.Origin origin, + long ifSeqNo, + long ifPrimaryTerm + ) { + throw new UnsupportedOperationException("DataFormatAwareReadOnlyEngine does not support deletes"); + } + // ---- IndexerLifecycleOperations (Tasks 4, 7) ---- @Override diff --git a/server/src/main/java/org/opensearch/index/translog/Translog.java b/server/src/main/java/org/opensearch/index/translog/Translog.java index 0f5bec7c3e54a..a7d8468796c33 100644 --- a/server/src/main/java/org/opensearch/index/translog/Translog.java +++ b/server/src/main/java/org/opensearch/index/translog/Translog.java @@ -1540,7 +1540,14 @@ public Source getSource() { } private void write(final StreamOutput out) throws IOException { - final int format = out.getVersion().onOrAfter(Version.V_2_0_0) ? SERIALIZATION_FORMAT : FORMAT_NO_VERSION_TYPE; + final int format; + if (out.getVersion().onOrAfter(Version.V_3_6_0)) { + format = SERIALIZATION_FORMAT; + } else if (out.getVersion().onOrAfter(Version.V_2_0_0)) { + format = FORMAT_NO_DOC_TYPE; + } else { + format = FORMAT_NO_VERSION_TYPE; + } out.writeVInt(format); if (format < FORMAT_NO_DOC_TYPE) { out.writeString(MapperService.SINGLE_MAPPING_NAME); diff --git a/server/src/main/java/org/opensearch/indices/pollingingest/MessageProcessorRunnable.java b/server/src/main/java/org/opensearch/indices/pollingingest/MessageProcessorRunnable.java index 90e2e49a5e2b6..92ba71fb31aa5 100644 --- a/server/src/main/java/org/opensearch/indices/pollingingest/MessageProcessorRunnable.java +++ b/server/src/main/java/org/opensearch/indices/pollingingest/MessageProcessorRunnable.java @@ -296,6 +296,8 @@ protected MessageOperation getOperation(ShardUpdateMessage shardUpdateMessage, M "Delete operation is missing ID. Skipping message." ); } else { + // TODO: routing is not available from the ingestion message; if pull-based + // ingestion adds routing support, thread it through here as well. operation = new Engine.Delete( id, new Term(IdFieldMapper.NAME, Uid.encodeId(id)), diff --git a/server/src/test/java/org/opensearch/index/translog/LocalTranslogTests.java b/server/src/test/java/org/opensearch/index/translog/LocalTranslogTests.java index 7328bc75624b8..acb74d00722ef 100644 --- a/server/src/test/java/org/opensearch/index/translog/LocalTranslogTests.java +++ b/server/src/test/java/org/opensearch/index/translog/LocalTranslogTests.java @@ -3612,6 +3612,7 @@ public void testTranslogOpSerialization() throws Exception { Translog.Index serializedIndex = (Translog.Index) Translog.Operation.readOperation(in); assertEquals(index, serializedIndex); + String deleteRouting = wireVersion.onOrAfter(Version.V_3_6_0) ? "custom-routing" : null; Engine.Delete eDelete = new Engine.Delete( doc.id(), newUid(doc), @@ -3622,7 +3623,8 @@ public void testTranslogOpSerialization() throws Exception { Origin.PRIMARY, 0, SequenceNumbers.UNASSIGNED_SEQ_NO, - 0 + 0, + deleteRouting ); Engine.DeleteResult eDeleteResult = new Engine.DeleteResult(2, randomPrimaryTerm, randomSeqNum, true); Translog.Delete delete = new Translog.Delete(eDelete, eDeleteResult); @@ -3634,6 +3636,76 @@ public void testTranslogOpSerialization() throws Exception { in.setVersion(wireVersion); Translog.Delete serializedDelete = (Translog.Delete) Translog.Operation.readOperation(in); assertEquals(delete, serializedDelete); + assertEquals(deleteRouting, serializedDelete.routing()); + } + + public void testDeleteRoutingBackwardCompatibility() throws Exception { + // Old version writes without routing, new version reads with routing=null + BytesStreamOutput out = new BytesStreamOutput(); + Version oldVersion = VersionUtils.randomVersionBetween( + random(), + Version.CURRENT.minimumCompatibilityVersion(), + VersionUtils.getPreviousVersion(Version.V_3_6_0) + ); + out.setVersion(oldVersion); + Translog.Delete deleteWithRouting = new Translog.Delete("doc-1", 1, 1, 1, "tenant-abc"); + Translog.Operation.writeOperation(out, deleteWithRouting); + + StreamInput in = out.bytes().streamInput(); + in.setVersion(oldVersion); + Translog.Delete deserialized = (Translog.Delete) Translog.Operation.readOperation(in); + assertEquals("doc-1", deserialized.id()); + assertEquals(1, deserialized.seqNo()); + assertNull("Routing must be null when written with old format", deserialized.routing()); + + // New version writes with routing, new version reads with routing preserved + out = new BytesStreamOutput(); + out.setVersion(Version.CURRENT); + Translog.Operation.writeOperation(out, deleteWithRouting); + + in = out.bytes().streamInput(); + in.setVersion(Version.CURRENT); + deserialized = (Translog.Delete) Translog.Operation.readOperation(in); + assertEquals("doc-1", deserialized.id()); + assertEquals("tenant-abc", deserialized.routing()); + } + + public void testDeleteRoutingTranslogRoundTrip() throws Exception { + // Write deletes with and without routing to a real translog, read them back + Translog.Location loc1 = translog.add(new Translog.Delete("doc-1", 0, primaryTerm.get(), 1, "tenant-routing")); + Translog.Location loc2 = translog.add(new Translog.Delete("doc-2", 1, primaryTerm.get(), 1)); + + Translog.Delete readBack1 = (Translog.Delete) translog.readOperation(loc1); + assertNotNull(readBack1); + assertEquals("doc-1", readBack1.id()); + assertEquals("tenant-routing", readBack1.routing()); + + Translog.Delete readBack2 = (Translog.Delete) translog.readOperation(loc2); + assertNotNull(readBack2); + assertEquals("doc-2", readBack2.id()); + assertNull(readBack2.routing()); + + // Verify via snapshot as well + translog.rollGeneration(); + try (Translog.Snapshot snapshot = translog.newSnapshot()) { + Translog.Operation op; + boolean foundRouted = false; + boolean foundUnrouted = false; + while ((op = snapshot.next()) != null) { + if (op instanceof Translog.Delete) { + Translog.Delete del = (Translog.Delete) op; + if ("doc-1".equals(del.id())) { + assertEquals("tenant-routing", del.routing()); + foundRouted = true; + } else if ("doc-2".equals(del.id())) { + assertNull(del.routing()); + foundUnrouted = true; + } + } + } + assertTrue("Should find delete with routing in snapshot", foundRouted); + assertTrue("Should find delete without routing in snapshot", foundUnrouted); + } } public void testRollGeneration() throws Exception { From 2588f237a1acd80cade36d28cb0244baa26ce90d Mon Sep 17 00:00:00 2001 From: Arpit Singla Date: Tue, 28 Jul 2026 23:35:47 +0530 Subject: [PATCH 3/4] fix(api): preserve backward compat for delete routing public APIs - Make IndexerEngineOperations.prepareDelete(routing) a default method delegating to the no-routing overload, fixing japicmp breakage - Add deprecated overloads for IndexShard.applyDeleteOperationOnPrimary and applyDeleteOperationOnReplica without routing parameter - Fix misleading TombstoneDocSupplier javadoc on default method - Use static import for hamcrest not() in LuceneChangesSnapshotTests Signed-off-by: Arpit Singla --- .../opensearch/index/engine/EngineConfig.java | 3 ++- .../engine/exec/IndexerEngineOperations.java | 6 +++-- .../opensearch/index/shard/IndexShard.java | 27 +++++++++++++++++++ .../engine/LuceneChangesSnapshotTests.java | 3 ++- 4 files changed, 35 insertions(+), 4 deletions(-) diff --git a/server/src/main/java/org/opensearch/index/engine/EngineConfig.java b/server/src/main/java/org/opensearch/index/engine/EngineConfig.java index d74ff6c8a0a5a..0324a651ce0e9 100644 --- a/server/src/main/java/org/opensearch/index/engine/EngineConfig.java +++ b/server/src/main/java/org/opensearch/index/engine/EngineConfig.java @@ -617,7 +617,8 @@ public interface TombstoneDocSupplier { ParsedDocument newDeleteTombstoneDoc(String id); /** - * Creates a tombstone document for a delete operation, preserving the routing value. + * Creates a tombstone document for a delete operation with routing. + * Default ignores routing for backward compatibility; override to preserve it. */ default ParsedDocument newDeleteTombstoneDoc(String id, String routing) { return newDeleteTombstoneDoc(id); diff --git a/server/src/main/java/org/opensearch/index/engine/exec/IndexerEngineOperations.java b/server/src/main/java/org/opensearch/index/engine/exec/IndexerEngineOperations.java index 3968c0dc69666..7faa87019ec29 100644 --- a/server/src/main/java/org/opensearch/index/engine/exec/IndexerEngineOperations.java +++ b/server/src/main/java/org/opensearch/index/engine/exec/IndexerEngineOperations.java @@ -108,7 +108,7 @@ Engine.Delete prepareDelete( long ifPrimaryTerm ); - Engine.Delete prepareDelete( + default Engine.Delete prepareDelete( String id, String routing, long seqNo, @@ -118,5 +118,7 @@ Engine.Delete prepareDelete( Engine.Operation.Origin origin, long ifSeqNo, long ifPrimaryTerm - ); + ) { + return prepareDelete(id, seqNo, primaryTerm, version, versionType, origin, ifSeqNo, ifPrimaryTerm); + } } diff --git a/server/src/main/java/org/opensearch/index/shard/IndexShard.java b/server/src/main/java/org/opensearch/index/shard/IndexShard.java index 7ba42cfd0d7d6..79f28d9f9ae1f 100644 --- a/server/src/main/java/org/opensearch/index/shard/IndexShard.java +++ b/server/src/main/java/org/opensearch/index/shard/IndexShard.java @@ -1444,6 +1444,20 @@ public Engine.DeleteResult applyDeleteOperationOnPrimary( ); } + /** + * @deprecated Use {@link #applyDeleteOperationOnPrimary(long, String, String, VersionType, long, long)} instead. + */ + @Deprecated + public Engine.DeleteResult applyDeleteOperationOnPrimary( + long version, + String id, + VersionType versionType, + long ifSeqNo, + long ifPrimaryTerm + ) throws IOException { + return applyDeleteOperationOnPrimary(version, id, null, versionType, ifSeqNo, ifPrimaryTerm); + } + public Engine.DeleteResult applyDeleteOperationOnReplica( long seqNo, long opPrimaryTerm, @@ -1481,6 +1495,19 @@ public Engine.DeleteResult applyDeleteOperationOnReplica( ); } + /** + * @deprecated Use {@link #applyDeleteOperationOnReplica(long, long, long, String, String)} instead. + */ + @Deprecated + public Engine.DeleteResult applyDeleteOperationOnReplica( + long seqNo, + long opPrimaryTerm, + long version, + String id + ) throws IOException { + return applyDeleteOperationOnReplica(seqNo, opPrimaryTerm, version, id, null); + } + private Engine.DeleteResult applyDeleteOperation( Indexer engine, long seqNo, diff --git a/server/src/test/java/org/opensearch/index/engine/LuceneChangesSnapshotTests.java b/server/src/test/java/org/opensearch/index/engine/LuceneChangesSnapshotTests.java index 06d4bf140d0b3..b241550a1b5d6 100644 --- a/server/src/test/java/org/opensearch/index/engine/LuceneChangesSnapshotTests.java +++ b/server/src/test/java/org/opensearch/index/engine/LuceneChangesSnapshotTests.java @@ -55,6 +55,7 @@ import static org.hamcrest.Matchers.containsString; import static org.hamcrest.Matchers.equalTo; +import static org.hamcrest.Matchers.not; public class LuceneChangesSnapshotTests extends EngineTestCase { private MapperService mapperService; @@ -393,7 +394,7 @@ public void testDeleteRoutingSerialization() throws Exception { // Verify toString includes routing assertThat(deleteWithRouting.toString(), containsString("routing=" + routingValue)); - assertThat(deleteNoRouting.toString(), org.hamcrest.Matchers.not(containsString("routing="))); + assertThat(deleteNoRouting.toString(), not(containsString("routing="))); } public void testOverFlow() throws Exception { From 802edd28402fca768afaf7927104b086b575b42a Mon Sep 17 00:00:00 2001 From: Arpit Singla Date: Tue, 28 Jul 2026 23:42:49 +0530 Subject: [PATCH 4/4] fix(translog): use V_3_8_0 version gate instead of V_3_6_0 V_3_6_0 and V_3_7_0 are already released without routing support. Using V_3_6_0 would cause TranslogCorruptedException on 3.6/3.7 nodes in mixed-version clusters. Gate on V_3_8_0 (Version.CURRENT) which is the first version that will ship with this feature. Signed-off-by: Arpit Singla --- .../src/main/java/org/opensearch/index/translog/Translog.java | 2 +- .../org/opensearch/index/translog/LocalTranslogTests.java | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/server/src/main/java/org/opensearch/index/translog/Translog.java b/server/src/main/java/org/opensearch/index/translog/Translog.java index a7d8468796c33..9edcff26b0112 100644 --- a/server/src/main/java/org/opensearch/index/translog/Translog.java +++ b/server/src/main/java/org/opensearch/index/translog/Translog.java @@ -1541,7 +1541,7 @@ public Source getSource() { private void write(final StreamOutput out) throws IOException { final int format; - if (out.getVersion().onOrAfter(Version.V_3_6_0)) { + if (out.getVersion().onOrAfter(Version.V_3_8_0)) { format = SERIALIZATION_FORMAT; } else if (out.getVersion().onOrAfter(Version.V_2_0_0)) { format = FORMAT_NO_DOC_TYPE; diff --git a/server/src/test/java/org/opensearch/index/translog/LocalTranslogTests.java b/server/src/test/java/org/opensearch/index/translog/LocalTranslogTests.java index acb74d00722ef..3258dd20f9f2b 100644 --- a/server/src/test/java/org/opensearch/index/translog/LocalTranslogTests.java +++ b/server/src/test/java/org/opensearch/index/translog/LocalTranslogTests.java @@ -3612,7 +3612,7 @@ public void testTranslogOpSerialization() throws Exception { Translog.Index serializedIndex = (Translog.Index) Translog.Operation.readOperation(in); assertEquals(index, serializedIndex); - String deleteRouting = wireVersion.onOrAfter(Version.V_3_6_0) ? "custom-routing" : null; + String deleteRouting = wireVersion.onOrAfter(Version.V_3_8_0) ? "custom-routing" : null; Engine.Delete eDelete = new Engine.Delete( doc.id(), newUid(doc), @@ -3645,7 +3645,7 @@ public void testDeleteRoutingBackwardCompatibility() throws Exception { Version oldVersion = VersionUtils.randomVersionBetween( random(), Version.CURRENT.minimumCompatibilityVersion(), - VersionUtils.getPreviousVersion(Version.V_3_6_0) + VersionUtils.getPreviousVersion(Version.V_3_8_0) ); out.setVersion(oldVersion); Translog.Delete deleteWithRouting = new Translog.Delete("doc-1", 1, 1, 1, "tenant-abc");