Skip to content
Merged
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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/),
### Changed

### Fixed
- Prevent criteria update for context aware indices ([#20250](https://github.com/opensearch-project/OpenSearch/pull/20250))

### Dependencies

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,13 +32,16 @@

package org.opensearch.common.lucene.uid;

import org.apache.lucene.index.FilterLeafReader;
import org.apache.lucene.index.IndexReader;
import org.apache.lucene.index.LeafReader;
import org.apache.lucene.index.LeafReaderContext;
import org.apache.lucene.index.SegmentReader;
import org.apache.lucene.index.Term;
import org.apache.lucene.util.CloseableThreadLocal;
import org.opensearch.common.annotation.PublicApi;
import org.opensearch.common.util.concurrent.ConcurrentCollections;
import org.opensearch.index.codec.CriteriaBasedCodec;

import java.io.IOException;
import java.util.List;
Expand Down Expand Up @@ -155,7 +158,8 @@ public static class DocIdAndSeqNo {
* <li>a doc ID and a version otherwise
* </ul>
*/
public static DocIdAndVersion loadDocIdAndVersion(IndexReader reader, Term term, boolean loadSeqNo) throws IOException {
public static DocIdAndVersion loadDocIdAndVersion(IndexReader reader, Term term, boolean loadSeqNo, String currentCriteria)
throws IOException {
PerThreadIDVersionAndSeqNoLookup[] lookups = getLookupState(reader, term.field());
List<LeafReaderContext> leaves = reader.leaves();
// iterate backwards to optimize for the frequently updated documents
Expand All @@ -165,6 +169,18 @@ public static DocIdAndVersion loadDocIdAndVersion(IndexReader reader, Term term,
PerThreadIDVersionAndSeqNoLookup lookup = lookups[leaf.ord];
DocIdAndVersion result = lookup.lookupVersion(term.bytes(), loadSeqNo, leaf);
if (result != null) {
if (result.version > 0 && currentCriteria != null) {
SegmentReader unwrappedReader = (SegmentReader) (FilterLeafReader.unwrap(leaf.reader()));
String prevCriteria = unwrappedReader.getSegmentInfo().info.getAttribute(CriteriaBasedCodec.BUCKET_NAME);
assert prevCriteria != null;
if (prevCriteria.equals(currentCriteria) == false) {
Comment thread
RS146BIJAY marked this conversation as resolved.
throw new UnsupportedOperationException(
"Updating grouping criteria is not allowed for context aware enabled indices.",
null
);
}
}

return result;
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,6 @@
import org.opensearch.common.logging.Loggers;
import org.opensearch.common.unit.TimeValue;
import org.opensearch.common.util.concurrent.ConcurrentCollections;
import org.opensearch.common.util.concurrent.KeyedLock;
import org.opensearch.common.util.concurrent.ReleasableLock;
import org.opensearch.common.util.io.IOUtils;
import org.opensearch.core.Assertions;
Expand Down Expand Up @@ -117,8 +116,6 @@
*/
public class CompositeIndexWriter implements DocumentIndexWriter {

private final KeyedLock<BytesRef> keyedLock = new KeyedLock<>();

private final EngineConfig engineConfig;
private final IndexWriter accumulatingIndexWriter;
private final CheckedBiFunction<String, CriteriaBasedIndexWriterLookup, DisposableIndexWriter, IOException> childIndexWriterFactory;
Expand Down Expand Up @@ -618,10 +615,6 @@ public void afterRefresh(boolean didRefresh) throws IOException {
liveIndexWriterDeletesMap = liveIndexWriterDeletesMap.invalidateOldMap();
}

Releasable acquireLock(BytesRef uid) {
return keyedLock.acquire(uid);
}

public Map<BytesRef, DeleteEntry> getLastDeleteEntrySet() {
return liveIndexWriterDeletesMap.old.lastDeleteEntrySet;
}
Expand All @@ -631,19 +624,16 @@ void putLastDeleteEntryUnderLockInNewMap(BytesRef uid, DeleteEntry entry) {
}

void putCriteria(BytesRef uid, String criteria) {
assert assertKeyedLockHeldByCurrentThread(uid);
assert uid.bytes.length == uid.length : "Oversized _uid! UID length: " + uid.length + ", bytes length: " + uid.bytes.length;
liveIndexWriterDeletesMap.putCriteriaForDoc(uid, criteria);
}

DisposableIndexWriter getIndexWriterForIdFromCurrent(BytesRef uid) {
assert assertKeyedLockHeldByCurrentThread(uid);
assert uid.bytes.length == uid.length : "Oversized _uid! UID length: " + uid.length + ", bytes length: " + uid.bytes.length;
return getIndexWriterForIdFromLookup(uid, liveIndexWriterDeletesMap.current);
}

DisposableIndexWriter getIndexWriterForIdFromOld(BytesRef uid) {
assert assertKeyedLockHeldByCurrentThread(uid);
assert uid.bytes.length == uid.length : "Oversized _uid! UID length: " + uid.length + ", bytes length: " + uid.bytes.length;
return getIndexWriterForIdFromLookup(uid, liveIndexWriterDeletesMap.old);
}
Expand Down Expand Up @@ -674,13 +664,21 @@ public boolean hasNewIndexingOrUpdates() {
return liveIndexWriterDeletesMap.hasNewIndexingOrUpdates();
}

String getCriteriaForDoc(BytesRef uid) {
return liveIndexWriterDeletesMap.getCriteriaForDoc(uid);
public boolean validateImmutableFieldNotUpdated(ParseContext.Document currentDocument, BytesRef currentUID) {
String currentCriteria = currentDocument.getGroupingCriteria();
String previousCriteria = getCriteriaForUID(currentUID, liveIndexWriterDeletesMap);
// previousCriteria may be null in case version is coming from a tombstone. In that case, we should ignore
// it.
return previousCriteria != null && previousCriteria.equals(currentCriteria) == false;
}

boolean assertKeyedLockHeldByCurrentThread(BytesRef uid) {
assert keyedLock.isHeldByCurrentThread(uid) : "Thread [" + Thread.currentThread().getName() + "], uid [" + uid.utf8ToString() + "]";
return true;
private String getCriteriaForUID(BytesRef uid, LiveIndexWriterDeletesMap currentMap) {
String criteria = currentMap.current.getCriteriaForDoc(uid);
if (criteria != null) {
return criteria;
}

return currentMap.old.getCriteriaForDoc(uid);
}

DisposableIndexWriter computeIndexWriterIfAbsentForCriteria(
Expand Down Expand Up @@ -925,10 +923,7 @@ public long addDocuments(final List<ParseContext.Document> docs, Term uid) throw
ensureOpen();
final String criteria = getGroupingCriteriaForDoc(docs.iterator().next());
DisposableIndexWriter disposableIndexWriter = getAssociatedIndexWriterForCriteria(criteria);
try (
CriteriaBasedIndexWriterLookup.CriteriaBasedWriterLock ignoreLock = disposableIndexWriter.getLookupMap().getMapReadLock();
Releasable ignore1 = acquireLock(uid.bytes())
) {
try (CriteriaBasedIndexWriterLookup.CriteriaBasedWriterLock ignoreLock = disposableIndexWriter.getLookupMap().getMapReadLock()) {
putCriteria(uid.bytes(), criteria);
long seqNo = disposableIndexWriter.getIndexWriter().addDocuments(docs);
childWriterPendingNumDocs.addAndGet(docs.size());
Expand All @@ -941,10 +936,7 @@ public long addDocument(ParseContext.Document doc, Term uid) throws IOException
ensureOpen();
final String criteria = getGroupingCriteriaForDoc(doc);
DisposableIndexWriter disposableIndexWriter = getAssociatedIndexWriterForCriteria(criteria);
try (
CriteriaBasedIndexWriterLookup.CriteriaBasedWriterLock ignoreLock = disposableIndexWriter.getLookupMap().getMapReadLock();
Releasable ignore1 = acquireLock(uid.bytes())
) {
try (CriteriaBasedIndexWriterLookup.CriteriaBasedWriterLock ignoreLock = disposableIndexWriter.getLookupMap().getMapReadLock()) {
putCriteria(uid.bytes(), criteria);
long seqNo = disposableIndexWriter.getIndexWriter().addDocument(doc);
childWriterPendingNumDocs.incrementAndGet();
Expand All @@ -964,10 +956,7 @@ public void softUpdateDocuments(
ensureOpen();
final String criteria = getGroupingCriteriaForDoc(docs.iterator().next());
DisposableIndexWriter disposableIndexWriter = getAssociatedIndexWriterForCriteria(criteria);
try (
CriteriaBasedIndexWriterLookup.CriteriaBasedWriterLock ignoreLock = disposableIndexWriter.getLookupMap().getMapReadLock();
Releasable ignore1 = acquireLock(uid.bytes())
) {
try (CriteriaBasedIndexWriterLookup.CriteriaBasedWriterLock ignoreLock = disposableIndexWriter.getLookupMap().getMapReadLock()) {
putCriteria(uid.bytes(), criteria);
disposableIndexWriter.getIndexWriter().softUpdateDocuments(uid, docs, softDeletesField);
childWriterPendingNumDocs.addAndGet(docs.size());
Expand All @@ -991,10 +980,7 @@ public void softUpdateDocument(
ensureOpen();
final String criteria = getGroupingCriteriaForDoc(doc);
DisposableIndexWriter disposableIndexWriter = getAssociatedIndexWriterForCriteria(criteria);
try (
CriteriaBasedIndexWriterLookup.CriteriaBasedWriterLock ignoreLock = disposableIndexWriter.getLookupMap().getMapReadLock();
Releasable ignore1 = acquireLock(uid.bytes())
) {
try (CriteriaBasedIndexWriterLookup.CriteriaBasedWriterLock ignoreLock = disposableIndexWriter.getLookupMap().getMapReadLock()) {
putCriteria(uid.bytes(), criteria);
disposableIndexWriter.getIndexWriter().softUpdateDocument(uid, doc, softDeletesField);
childWriterPendingNumDocs.incrementAndGet();
Expand Down Expand Up @@ -1031,33 +1017,29 @@ public void deleteDocument(
Field... softDeletesField
) throws IOException {
ensureOpen();
try (Releasable ignore1 = acquireLock(uid.bytes())) {
CompositeIndexWriter.DisposableIndexWriter currentDisposableWriter = getIndexWriterForIdFromCurrent(uid.bytes());
if (currentDisposableWriter != null) {
try (
CriteriaBasedIndexWriterLookup.CriteriaBasedWriterLock ignore = currentDisposableWriter.getLookupMap().getMapReadLock()
) {
if (currentDisposableWriter.getLookupMap().isClosed() == false && isStaleOperation == false) {
addDeleteEntryToWriter(new DeleteEntry(uid, version, seqNo, primaryTerm), currentDisposableWriter.getIndexWriter());
// only increment this when addDeleteEntry for child writers are called.
childWriterPendingNumDocs.incrementAndGet();
}
CompositeIndexWriter.DisposableIndexWriter currentDisposableWriter = getIndexWriterForIdFromCurrent(uid.bytes());
if (currentDisposableWriter != null) {
try (CriteriaBasedIndexWriterLookup.CriteriaBasedWriterLock ignore = currentDisposableWriter.getLookupMap().getMapReadLock()) {
if (currentDisposableWriter.getLookupMap().isClosed() == false && isStaleOperation == false) {
addDeleteEntryToWriter(new DeleteEntry(uid, version, seqNo, primaryTerm), currentDisposableWriter.getIndexWriter());
// only increment this when addDeleteEntry for child writers are called.
childWriterPendingNumDocs.incrementAndGet();
}
}
}

CompositeIndexWriter.DisposableIndexWriter oldDisposableWriter = getIndexWriterForIdFromOld(uid.bytes());
if (oldDisposableWriter != null) {
try (CriteriaBasedIndexWriterLookup.CriteriaBasedWriterLock ignore = oldDisposableWriter.getLookupMap().getMapReadLock()) {
if (oldDisposableWriter.getLookupMap().isClosed() == false && isStaleOperation == false) {
addDeleteEntryToWriter(new DeleteEntry(uid, version, seqNo, primaryTerm), oldDisposableWriter.getIndexWriter());
// only increment this when addDeleteEntry for child writers are called.
childWriterPendingNumDocs.incrementAndGet();
}
CompositeIndexWriter.DisposableIndexWriter oldDisposableWriter = getIndexWriterForIdFromOld(uid.bytes());
if (oldDisposableWriter != null) {
try (CriteriaBasedIndexWriterLookup.CriteriaBasedWriterLock ignore = oldDisposableWriter.getLookupMap().getMapReadLock()) {
if (oldDisposableWriter.getLookupMap().isClosed() == false && isStaleOperation == false) {
addDeleteEntryToWriter(new DeleteEntry(uid, version, seqNo, primaryTerm), oldDisposableWriter.getIndexWriter());
// only increment this when addDeleteEntry for child writers are called.
childWriterPendingNumDocs.incrementAndGet();
}
}

deleteInLucene(uid, isStaleOperation, accumulatingIndexWriter, doc, softDeletesField);
}

deleteInLucene(uid, isStaleOperation, accumulatingIndexWriter, doc, softDeletesField);
}

private void deleteInLucene(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
import org.apache.lucene.index.LiveIndexWriterConfig;
import org.apache.lucene.index.Term;
import org.apache.lucene.search.ReferenceManager;
import org.apache.lucene.util.BytesRef;
import org.opensearch.common.lease.Releasable;
import org.opensearch.index.mapper.ParseContext;

Expand Down Expand Up @@ -92,4 +93,6 @@ void deleteDocument(
boolean isWriteLockedByCurrentThread();

Releasable obtainWriteLockOnAllMap();

boolean validateImmutableFieldNotUpdated(ParseContext.Document previousDocument, BytesRef currentUID);
}
Original file line number Diff line number Diff line change
Expand Up @@ -335,7 +335,8 @@ protected long getMaxSeqNoFromSearcher(IndexSearcher searcher) throws IOExceptio
VersionsAndSeqNoResolver.DocIdAndVersion docIdAndVersion = VersionsAndSeqNoResolver.loadDocIdAndVersion(
searcher.getIndexReader(),
uidTerm,
true
true,
null
);
assert docIdAndVersion != null;
return docIdAndVersion.seqNo;
Expand Down Expand Up @@ -704,7 +705,7 @@ protected final GetResult getFromSearcher(
final Engine.Searcher searcher = searcherFactory.apply("get", scope);
final DocIdAndVersion docIdAndVersion;
try {
docIdAndVersion = VersionsAndSeqNoResolver.loadDocIdAndVersion(searcher.getIndexReader(), get.uid(), true);
docIdAndVersion = VersionsAndSeqNoResolver.loadDocIdAndVersion(searcher.getIndexReader(), get.uid(), true, null);
} catch (Exception e) {
Releasables.closeWhileHandlingException(searcher);
// TODO: A better exception goes here
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -781,7 +781,19 @@ protected VersionValue resolveDocVersion(final Operation op, boolean loadSeqNo)
// IndexWriters with parent writers, version will be either present in version map or in parent IndexWriter. So we do not need
// to resolve version from child level IndexWriters (both from mark for refresh and active IndexWriter).
try (Searcher searcher = acquireSearcher("load_version", SearcherScope.INTERNAL)) {
docIdAndVersion = VersionsAndSeqNoResolver.loadDocIdAndVersion(searcher.getIndexReader(), op.uid(), loadSeqNo);
String currentCriteria = null;
if (op instanceof Index) {
Index index = (Index) op;
assert index.docs() != null && index.docs().isEmpty() == false;
currentCriteria = index.docs().get(0).getGroupingCriteria();
}

docIdAndVersion = VersionsAndSeqNoResolver.loadDocIdAndVersion(
searcher.getIndexReader(),
op.uid(),
loadSeqNo,
currentCriteria
);
}
if (docIdAndVersion != null) {
versionValue = new IndexVersionValue(null, docIdAndVersion.version, docIdAndVersion.seqNo, docIdAndVersion.primaryTerm);
Expand All @@ -790,6 +802,15 @@ protected VersionValue resolveDocVersion(final Operation op, boolean loadSeqNo)
&& versionValue.isDelete()
&& (engineConfig.getThreadPool().relativeTimeInMillis() - ((DeleteVersionValue) versionValue).time) > getGcDeletesInMillis()) {
versionValue = null;
} else if (op instanceof Index && versionValue.version > 0) {
Index index = (Index) op;
assert index.docs() != null && index.docs().isEmpty() == false;
if (documentIndexWriter.validateImmutableFieldNotUpdated(index.docs().get(0), op.uid().bytes())) {
throw new UnsupportedOperationException(
"Updating grouping criteria is not allowed for context aware enabled indices.",
null
);
}
}
return versionValue;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@
import org.apache.lucene.index.IndexWriter;
import org.apache.lucene.index.LiveIndexWriterConfig;
import org.apache.lucene.index.Term;
import org.apache.lucene.util.BytesRef;
import org.opensearch.common.lease.Releasable;
import org.opensearch.index.mapper.ParseContext;

Expand Down Expand Up @@ -226,4 +227,8 @@ public void afterRefresh(boolean b) throws IOException {
public Releasable obtainWriteLockOnAllMap() {
return () -> {};
}

public boolean validateImmutableFieldNotUpdated(ParseContext.Document previousDocument, BytesRef currentUID) {
return false;
}
}
Loading
Loading