Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -653,6 +653,7 @@ static boolean executeBulkItemRequest(
result = primary.applyDeleteOperationOnPrimary(
version,
request.id(),
request.routing(),
request.versionType(),
request.ifSeqNo(),
request.ifPrimaryTerm()
Expand Down Expand Up @@ -935,7 +936,8 @@ private static Engine.Result performOpOnReplica(
primaryResponse.getSeqNo(),
primaryResponse.getPrimaryTerm(),
primaryResponse.getVersion(),
deleteRequest.id()
deleteRequest.id(),
deleteRequest.routing()
);
break;
default:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
* <p>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
47 changes: 44 additions & 3 deletions server/src/main/java/org/opensearch/index/engine/Engine.java
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}

/**
Expand Down Expand Up @@ -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;

Expand All @@ -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;
Expand All @@ -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;
}
Expand All @@ -1944,7 +1976,8 @@ public Delete(String id, Term uid, long primaryTerm) {
Origin.PRIMARY,
System.nanoTime(),
UNASSIGNED_SEQ_NO,
0
0,
null
);
}

Expand All @@ -1959,7 +1992,8 @@ public Delete(Delete template, VersionType versionType) {
template.origin(),
template.startTime(),
UNASSIGNED_SEQ_NO,
0
0,
template.routing()
);
}

Expand All @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<IndexCommit> acquireSafeIndexCommit() throws EngineException {
return engine.acquireSafeIndexCommit();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -208,6 +208,7 @@ public DocumentMapper(MapperService mapperService, Mapping mapping, long version
final Collection<String> deleteTombstoneMetadataFields = Arrays.asList(
VersionFieldMapper.NAME,
IdFieldMapper.NAME,
RoutingFieldMapper.NAME,
SeqNoFieldMapper.NAME,
SeqNoFieldMapper.PRIMARY_TERM_NAME,
SeqNoFieldMapper.TOMBSTONE_NAME
Expand Down Expand Up @@ -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();
}

Expand Down
Loading
Loading