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/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/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..0324a651ce0e9 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,14 @@ public interface TombstoneDocSupplier { */ ParsedDocument newDeleteTombstoneDoc(String id); + /** + * 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); + } + /** * 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..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 @@ -107,4 +107,18 @@ Engine.Delete prepareDelete( long ifSeqNo, long ifPrimaryTerm ); + + default Engine.Delete prepareDelete( + String id, + String routing, + long seqNo, + long primaryTerm, + long version, + VersionType versionType, + 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/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..79f28d9f9ae1f 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,27 @@ public Engine.DeleteResult applyDeleteOperationOnPrimary( ); } - public Engine.DeleteResult applyDeleteOperationOnReplica(long seqNo, long opPrimaryTerm, long version, String id) throws IOException { + /** + * @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, + long version, + String id, + @Nullable String routing + ) throws IOException { if (indexSettings.isSegRepEnabledOrRemoteNode()) { final Engine.Delete delete = new Engine.Delete( id, @@ -1454,7 +1476,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 +1487,7 @@ public Engine.DeleteResult applyDeleteOperationOnReplica(long seqNo, long opPrim opPrimaryTerm, version, id, + routing, null, UNASSIGNED_SEQ_NO, 0, @@ -1471,12 +1495,26 @@ public Engine.DeleteResult applyDeleteOperationOnReplica(long seqNo, long opPrim ); } + /** + * @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, long opPrimaryTerm, long version, String id, + @Nullable String routing, @Nullable VersionType versionType, long ifSeqNo, long ifPrimaryTerm, @@ -1488,7 +1526,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 +1545,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 +3080,7 @@ private Engine.Result applyTranslogOperation(Indexer engine, Translog.Operation delete.primaryTerm(), delete.version(), delete.id(), + delete.routing(), versionType, UNASSIGNED_SEQ_NO, 0, @@ -5645,6 +5699,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..9edcff26b0112 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; @@ -1521,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_8_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); @@ -1537,6 +1563,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 +1579,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 +1588,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/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/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..b241550a1b5d6 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; @@ -53,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; @@ -358,6 +361,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(), 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..3258dd20f9f2b 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)); @@ -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_8_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_8_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 { 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;