diff --git a/CHANGELOG.md b/CHANGELOG.md index 9addad8988e29..cefab92355b26 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/server/src/main/java/org/opensearch/common/lucene/uid/VersionsAndSeqNoResolver.java b/server/src/main/java/org/opensearch/common/lucene/uid/VersionsAndSeqNoResolver.java index 4958267b1775b..11b1524dc2a15 100644 --- a/server/src/main/java/org/opensearch/common/lucene/uid/VersionsAndSeqNoResolver.java +++ b/server/src/main/java/org/opensearch/common/lucene/uid/VersionsAndSeqNoResolver.java @@ -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; @@ -155,7 +158,8 @@ public static class DocIdAndSeqNo { *
  • a doc ID and a version otherwise * */ - 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 leaves = reader.leaves(); // iterate backwards to optimize for the frequently updated documents @@ -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) { + throw new UnsupportedOperationException( + "Updating grouping criteria is not allowed for context aware enabled indices.", + null + ); + } + } + return result; } } diff --git a/server/src/main/java/org/opensearch/index/engine/CompositeIndexWriter.java b/server/src/main/java/org/opensearch/index/engine/CompositeIndexWriter.java index 96c8bb98c6d4d..5c421933529eb 100644 --- a/server/src/main/java/org/opensearch/index/engine/CompositeIndexWriter.java +++ b/server/src/main/java/org/opensearch/index/engine/CompositeIndexWriter.java @@ -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; @@ -117,8 +116,6 @@ */ public class CompositeIndexWriter implements DocumentIndexWriter { - private final KeyedLock keyedLock = new KeyedLock<>(); - private final EngineConfig engineConfig; private final IndexWriter accumulatingIndexWriter; private final CheckedBiFunction childIndexWriterFactory; @@ -618,10 +615,6 @@ public void afterRefresh(boolean didRefresh) throws IOException { liveIndexWriterDeletesMap = liveIndexWriterDeletesMap.invalidateOldMap(); } - Releasable acquireLock(BytesRef uid) { - return keyedLock.acquire(uid); - } - public Map getLastDeleteEntrySet() { return liveIndexWriterDeletesMap.old.lastDeleteEntrySet; } @@ -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); } @@ -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( @@ -925,10 +923,7 @@ public long addDocuments(final List 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()); @@ -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(); @@ -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()); @@ -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(); @@ -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( diff --git a/server/src/main/java/org/opensearch/index/engine/DocumentIndexWriter.java b/server/src/main/java/org/opensearch/index/engine/DocumentIndexWriter.java index 3bfbd057fd720..607a24b33c816 100644 --- a/server/src/main/java/org/opensearch/index/engine/DocumentIndexWriter.java +++ b/server/src/main/java/org/opensearch/index/engine/DocumentIndexWriter.java @@ -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; @@ -92,4 +93,6 @@ void deleteDocument( boolean isWriteLockedByCurrentThread(); Releasable obtainWriteLockOnAllMap(); + + boolean validateImmutableFieldNotUpdated(ParseContext.Document previousDocument, BytesRef currentUID); } 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 eff9809667a3a..3ec1788a70309 100644 --- a/server/src/main/java/org/opensearch/index/engine/Engine.java +++ b/server/src/main/java/org/opensearch/index/engine/Engine.java @@ -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; @@ -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 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 a583e41213d1c..95cf0a05d1fa0 100644 --- a/server/src/main/java/org/opensearch/index/engine/InternalEngine.java +++ b/server/src/main/java/org/opensearch/index/engine/InternalEngine.java @@ -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); @@ -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; } diff --git a/server/src/main/java/org/opensearch/index/engine/LuceneIndexWriter.java b/server/src/main/java/org/opensearch/index/engine/LuceneIndexWriter.java index 689b61fd56ee2..5f165e51ae8ed 100644 --- a/server/src/main/java/org/opensearch/index/engine/LuceneIndexWriter.java +++ b/server/src/main/java/org/opensearch/index/engine/LuceneIndexWriter.java @@ -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; @@ -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; + } } diff --git a/server/src/test/java/org/opensearch/common/lucene/uid/VersionsTests.java b/server/src/test/java/org/opensearch/common/lucene/uid/VersionsTests.java index acefc5720bdfd..2ca9f74946a1c 100644 --- a/server/src/test/java/org/opensearch/common/lucene/uid/VersionsTests.java +++ b/server/src/test/java/org/opensearch/common/lucene/uid/VersionsTests.java @@ -77,7 +77,7 @@ public void testVersions() throws Exception { Directory dir = newDirectory(); IndexWriter writer = new IndexWriter(dir, new IndexWriterConfig(Lucene.STANDARD_ANALYZER)); DirectoryReader directoryReader = OpenSearchDirectoryReader.wrap(DirectoryReader.open(writer), new ShardId("foo", "_na_", 1)); - assertThat(loadDocIdAndVersion(directoryReader, new Term(IdFieldMapper.NAME, "1"), randomBoolean()), nullValue()); + assertThat(loadDocIdAndVersion(directoryReader, new Term(IdFieldMapper.NAME, "1"), randomBoolean(), null), nullValue()); Document doc = new Document(); doc.add(new Field(IdFieldMapper.NAME, "1", IdFieldMapper.Defaults.FIELD_TYPE)); @@ -86,7 +86,7 @@ public void testVersions() throws Exception { doc.add(new NumericDocValuesField(SeqNoFieldMapper.PRIMARY_TERM_NAME, randomLongBetween(1, Long.MAX_VALUE))); writer.updateDocument(new Term(IdFieldMapper.NAME, "1"), doc); directoryReader = reopen(directoryReader); - assertThat(loadDocIdAndVersion(directoryReader, new Term(IdFieldMapper.NAME, "1"), randomBoolean()).version, equalTo(1L)); + assertThat(loadDocIdAndVersion(directoryReader, new Term(IdFieldMapper.NAME, "1"), randomBoolean(), null).version, equalTo(1L)); doc = new Document(); Field uid = new Field(IdFieldMapper.NAME, "1", IdFieldMapper.Defaults.FIELD_TYPE); @@ -97,7 +97,7 @@ public void testVersions() throws Exception { doc.add(new NumericDocValuesField(SeqNoFieldMapper.PRIMARY_TERM_NAME, randomLongBetween(1, Long.MAX_VALUE))); writer.updateDocument(new Term(IdFieldMapper.NAME, "1"), doc); directoryReader = reopen(directoryReader); - assertThat(loadDocIdAndVersion(directoryReader, new Term(IdFieldMapper.NAME, "1"), randomBoolean()).version, equalTo(2L)); + assertThat(loadDocIdAndVersion(directoryReader, new Term(IdFieldMapper.NAME, "1"), randomBoolean(), null).version, equalTo(2L)); // test reuse of uid field doc = new Document(); @@ -109,11 +109,11 @@ public void testVersions() throws Exception { writer.updateDocument(new Term(IdFieldMapper.NAME, "1"), doc); directoryReader = reopen(directoryReader); - assertThat(loadDocIdAndVersion(directoryReader, new Term(IdFieldMapper.NAME, "1"), randomBoolean()).version, equalTo(3L)); + assertThat(loadDocIdAndVersion(directoryReader, new Term(IdFieldMapper.NAME, "1"), randomBoolean(), null).version, equalTo(3L)); writer.deleteDocuments(new Term(IdFieldMapper.NAME, "1")); directoryReader = reopen(directoryReader); - assertThat(loadDocIdAndVersion(directoryReader, new Term(IdFieldMapper.NAME, "1"), randomBoolean()), nullValue()); + assertThat(loadDocIdAndVersion(directoryReader, new Term(IdFieldMapper.NAME, "1"), randomBoolean(), null), nullValue()); directoryReader.close(); writer.close(); dir.close(); @@ -141,18 +141,18 @@ public void testNestedDocuments() throws IOException { writer.updateDocuments(new Term(IdFieldMapper.NAME, "1"), docs); DirectoryReader directoryReader = OpenSearchDirectoryReader.wrap(DirectoryReader.open(writer), new ShardId("foo", "_na_", 1)); - assertThat(loadDocIdAndVersion(directoryReader, new Term(IdFieldMapper.NAME, "1"), randomBoolean()).version, equalTo(5L)); + assertThat(loadDocIdAndVersion(directoryReader, new Term(IdFieldMapper.NAME, "1"), randomBoolean(), null).version, equalTo(5L)); version.setLongValue(6L); writer.updateDocuments(new Term(IdFieldMapper.NAME, "1"), docs); version.setLongValue(7L); writer.updateDocuments(new Term(IdFieldMapper.NAME, "1"), docs); directoryReader = reopen(directoryReader); - assertThat(loadDocIdAndVersion(directoryReader, new Term(IdFieldMapper.NAME, "1"), randomBoolean()).version, equalTo(7L)); + assertThat(loadDocIdAndVersion(directoryReader, new Term(IdFieldMapper.NAME, "1"), randomBoolean(), null).version, equalTo(7L)); writer.deleteDocuments(new Term(IdFieldMapper.NAME, "1")); directoryReader = reopen(directoryReader); - assertThat(loadDocIdAndVersion(directoryReader, new Term(IdFieldMapper.NAME, "1"), randomBoolean()), nullValue()); + assertThat(loadDocIdAndVersion(directoryReader, new Term(IdFieldMapper.NAME, "1"), randomBoolean(), null), nullValue()); directoryReader.close(); writer.close(); dir.close(); @@ -172,10 +172,10 @@ public void testCache() throws Exception { writer.addDocument(doc); DirectoryReader reader = DirectoryReader.open(writer); // should increase cache size by 1 - assertEquals(87, loadDocIdAndVersion(reader, new Term(IdFieldMapper.NAME, "6"), randomBoolean()).version); + assertEquals(87, loadDocIdAndVersion(reader, new Term(IdFieldMapper.NAME, "6"), randomBoolean(), null).version); assertEquals(size + 1, VersionsAndSeqNoResolver.lookupStates.size()); // should be cache hit - assertEquals(87, loadDocIdAndVersion(reader, new Term(IdFieldMapper.NAME, "6"), randomBoolean()).version); + assertEquals(87, loadDocIdAndVersion(reader, new Term(IdFieldMapper.NAME, "6"), randomBoolean(), null).version); assertEquals(size + 1, VersionsAndSeqNoResolver.lookupStates.size()); reader.close(); @@ -198,11 +198,11 @@ public void testCacheFilterReader() throws Exception { doc.add(new NumericDocValuesField(SeqNoFieldMapper.PRIMARY_TERM_NAME, randomLongBetween(1, Long.MAX_VALUE))); writer.addDocument(doc); DirectoryReader reader = DirectoryReader.open(writer); - assertEquals(87, loadDocIdAndVersion(reader, new Term(IdFieldMapper.NAME, "6"), randomBoolean()).version); + assertEquals(87, loadDocIdAndVersion(reader, new Term(IdFieldMapper.NAME, "6"), randomBoolean(), null).version); assertEquals(size + 1, VersionsAndSeqNoResolver.lookupStates.size()); // now wrap the reader DirectoryReader wrapped = OpenSearchDirectoryReader.wrap(reader, new ShardId("bogus", "_na_", 5)); - assertEquals(87, loadDocIdAndVersion(wrapped, new Term(IdFieldMapper.NAME, "6"), randomBoolean()).version); + assertEquals(87, loadDocIdAndVersion(wrapped, new Term(IdFieldMapper.NAME, "6"), randomBoolean(), null).version); // same size map: core cache key is shared assertEquals(size + 1, VersionsAndSeqNoResolver.lookupStates.size()); diff --git a/server/src/test/java/org/opensearch/index/engine/CompositeIndexWriterForAppendTests.java b/server/src/test/java/org/opensearch/index/engine/CompositeIndexWriterForAppendTests.java index 8e50a382adfca..725d2bdb13425 100644 --- a/server/src/test/java/org/opensearch/index/engine/CompositeIndexWriterForAppendTests.java +++ b/server/src/test/java/org/opensearch/index/engine/CompositeIndexWriterForAppendTests.java @@ -180,14 +180,10 @@ public void testChildDirectoryDeletedPostRefresh() throws IOException, Interrupt ); Engine.Index operation = indexForDoc(createParsedDoc("id", null, DEFAULT_CRITERIA)); - try (Releasable ignore1 = compositeIndexWriter.acquireLock(operation.uid().bytes())) { - compositeIndexWriter.addDocuments(operation.docs(), operation.uid()); - } + compositeIndexWriter.addDocuments(operation.docs(), operation.uid()); operation = indexForDoc(createParsedDoc("id2", null, "testingNewCriteria")); - try (Releasable ignore1 = compositeIndexWriter.acquireLock(operation.uid().bytes())) { - compositeIndexWriter.addDocuments(operation.docs(), operation.uid()); - } + compositeIndexWriter.addDocuments(operation.docs(), operation.uid()); compositeIndexWriter.beforeRefresh(); compositeIndexWriter.afterRefresh(true); @@ -348,7 +344,7 @@ public long addDocuments(Iterable> String id = Integer.toString(randomIntBetween(1, 100)); Engine.Index operation = indexForDoc(createParsedDoc(id, null, DEFAULT_CRITERIA)); - try (Releasable ignore1 = compositeIndexWriter.acquireLock(operation.uid().bytes())) { + try { addDocException.set(new IOException("simulated")); expectThrows(IOException.class, () -> compositeIndexWriter.addDocuments(operation.docs(), operation.uid())); } finally { diff --git a/server/src/test/java/org/opensearch/index/engine/CompositeIndexWriterForUpdateAndDeletesTests.java b/server/src/test/java/org/opensearch/index/engine/CompositeIndexWriterForUpdateAndDeletesTests.java index 14966d086ec33..5fbee39e26fbf 100644 --- a/server/src/test/java/org/opensearch/index/engine/CompositeIndexWriterForUpdateAndDeletesTests.java +++ b/server/src/test/java/org/opensearch/index/engine/CompositeIndexWriterForUpdateAndDeletesTests.java @@ -9,7 +9,6 @@ package org.opensearch.index.engine; import org.apache.lucene.index.DirectoryReader; -import org.opensearch.common.lease.Releasable; import org.opensearch.common.util.io.IOUtils; import java.io.IOException; @@ -30,23 +29,19 @@ public void testDeleteWithDocumentInParentWriter() throws IOException { indexWriterFactory ); Engine.Index operation = indexForDoc(createParsedDoc(id, null, DEFAULT_CRITERIA)); - try (Releasable ignore1 = compositeIndexWriter.acquireLock(operation.uid().bytes())) { - compositeIndexWriter.addDocuments(operation.docs(), operation.uid()); - } + compositeIndexWriter.addDocuments(operation.docs(), operation.uid()); compositeIndexWriter.beforeRefresh(); compositeIndexWriter.afterRefresh(true); - try (Releasable ignore1 = compositeIndexWriter.acquireLock(operation.uid().bytes())) { - compositeIndexWriter.deleteDocument( - operation.uid(), - false, - newDeleteTombstoneDoc(id), - 1, - 2, - primaryTerm.get(), - softDeletesField - ); - } + compositeIndexWriter.deleteDocument( + operation.uid(), + false, + newDeleteTombstoneDoc(id), + 1, + 2, + primaryTerm.get(), + softDeletesField + ); compositeIndexWriter.beforeRefresh(); compositeIndexWriter.afterRefresh(true); @@ -72,18 +67,16 @@ public void testDeleteWithDocumentInChildWriter() throws IOException { indexWriterFactory ); Engine.Index operation = indexForDoc(createParsedDoc(id, null, DEFAULT_CRITERIA)); - try (Releasable ignore1 = compositeIndexWriter.acquireLock(operation.uid().bytes())) { - compositeIndexWriter.addDocuments(operation.docs(), operation.uid()); - compositeIndexWriter.deleteDocument( - operation.uid(), - false, - newDeleteTombstoneDoc(id), - 1, - 2, - primaryTerm.get(), - softDeletesField - ); - } + compositeIndexWriter.addDocuments(operation.docs(), operation.uid()); + compositeIndexWriter.deleteDocument( + operation.uid(), + false, + newDeleteTombstoneDoc(id), + 1, + 2, + primaryTerm.get(), + softDeletesField + ); compositeIndexWriter.beforeRefresh(); compositeIndexWriter.afterRefresh(true); @@ -110,26 +103,22 @@ public void testDeleteWithDocumentInBothChildAndParentWriter() throws IOExceptio indexWriterFactory ); Engine.Index operation = indexForDoc(createParsedDoc(id, null, DEFAULT_CRITERIA)); - try (Releasable ignore1 = compositeIndexWriter.acquireLock(operation.uid().bytes())) { - compositeIndexWriter.addDocuments(operation.docs(), operation.uid()); - } + compositeIndexWriter.addDocuments(operation.docs(), operation.uid()); compositeIndexWriter.beforeRefresh(); compositeIndexWriter.afterRefresh(true); operation = indexForDoc(createParsedDoc(id, null, DEFAULT_CRITERIA)); - try (Releasable ignore1 = compositeIndexWriter.acquireLock(operation.uid().bytes())) { - compositeIndexWriter.softUpdateDocuments(operation.uid(), operation.docs(), 2, 2, primaryTerm.get(), softDeletesField); - compositeIndexWriter.deleteDocument( - operation.uid(), - false, - newDeleteTombstoneDoc(id), - 1, - 2, - primaryTerm.get(), - softDeletesField - ); - } + compositeIndexWriter.softUpdateDocuments(operation.uid(), operation.docs(), 2, 2, primaryTerm.get(), softDeletesField); + compositeIndexWriter.deleteDocument( + operation.uid(), + false, + newDeleteTombstoneDoc(id), + 1, + 2, + primaryTerm.get(), + softDeletesField + ); compositeIndexWriter.beforeRefresh(); compositeIndexWriter.afterRefresh(true); @@ -154,9 +143,7 @@ public void testDeleteWithDocumentInOldChildWriter() throws IOException, Interru ); Engine.Index operation = indexForDoc(createParsedDoc(id, null, DEFAULT_CRITERIA)); - try (Releasable ignore1 = compositeIndexWriter.acquireLock(operation.uid().bytes())) { - compositeIndexWriter.addDocuments(operation.docs(), operation.uid()); - } + compositeIndexWriter.addDocuments(operation.docs(), operation.uid()); CompositeIndexWriter.CriteriaBasedIndexWriterLookup lock = compositeIndexWriter.acquireNewReadLock(); CountDownLatch latch = new CountDownLatch(1); @@ -207,17 +194,13 @@ public void testUpdateWithDocumentInParentIndexWriter() throws IOException { indexWriterFactory ); Engine.Index operation = indexForDoc(createParsedDoc(id, null, DEFAULT_CRITERIA)); - try (Releasable ignore1 = compositeIndexWriter.acquireLock(operation.uid().bytes())) { - compositeIndexWriter.addDocuments(operation.docs(), operation.uid()); - } + compositeIndexWriter.addDocuments(operation.docs(), operation.uid()); compositeIndexWriter.beforeRefresh(); compositeIndexWriter.afterRefresh(true); operation = indexForDoc(createParsedDoc(id, null, DEFAULT_CRITERIA)); - try (Releasable ignore1 = compositeIndexWriter.acquireLock(operation.uid().bytes())) { - compositeIndexWriter.softUpdateDocuments(operation.uid(), operation.docs(), 2, 2, primaryTerm.get(), softDeletesField); - } + compositeIndexWriter.softUpdateDocuments(operation.uid(), operation.docs(), 2, 2, primaryTerm.get(), softDeletesField); compositeIndexWriter.beforeRefresh(); compositeIndexWriter.afterRefresh(true); @@ -243,15 +226,10 @@ public void testUpdateWithDocumentInChildIndexWriter() throws IOException { indexWriterFactory ); Engine.Index operation = indexForDoc(createParsedDoc(id, null, DEFAULT_CRITERIA)); - try (Releasable ignore1 = compositeIndexWriter.acquireLock(operation.uid().bytes())) { - compositeIndexWriter.addDocuments(operation.docs(), operation.uid()); - } + compositeIndexWriter.addDocuments(operation.docs(), operation.uid()); operation = indexForDoc(createParsedDoc(id, null, DEFAULT_CRITERIA)); - try (Releasable ignore1 = compositeIndexWriter.acquireLock(operation.uid().bytes())) { - compositeIndexWriter.softUpdateDocuments(operation.uid(), operation.docs(), 2, 2, primaryTerm.get(), softDeletesField); - } - + compositeIndexWriter.softUpdateDocuments(operation.uid(), operation.docs(), 2, 2, primaryTerm.get(), softDeletesField); compositeIndexWriter.beforeRefresh(); compositeIndexWriter.afterRefresh(true); try (DirectoryReader directoryReader = DirectoryReader.open(compositeIndexWriter.getAccumulatingIndexWriter())) { diff --git a/server/src/test/java/org/opensearch/index/engine/InternalEngineTests.java b/server/src/test/java/org/opensearch/index/engine/InternalEngineTests.java index 34d0d52df6d74..bb6ad2251f9c7 100644 --- a/server/src/test/java/org/opensearch/index/engine/InternalEngineTests.java +++ b/server/src/test/java/org/opensearch/index/engine/InternalEngineTests.java @@ -1878,7 +1878,7 @@ public void testLookupVersionWithPrunedAwayIds() throws IOException { writer.forceMerge(1); try (DirectoryReader reader = DirectoryReader.open(writer)) { assertEquals(1, reader.leaves().size()); - assertNull(VersionsAndSeqNoResolver.loadDocIdAndVersion(reader, new Term(IdFieldMapper.NAME, "1"), false)); + assertNull(VersionsAndSeqNoResolver.loadDocIdAndVersion(reader, new Term(IdFieldMapper.NAME, "1"), false, null)); } } } @@ -8495,7 +8495,7 @@ public void testNewChangesSnapshotWithDeleteAndUpdateWithDerivedSourceAndContext ParsedDocument doc = testParsedDocument( Integer.toString(i), null, - testContextSpecificDocument(), + testContextSpecificDocument("grouping_criteria"), null, // No source, it should be derived null ); @@ -8686,7 +8686,7 @@ public CacheHelper getReaderCacheHelper() { ParsedDocument doc = testParsedDocument( Integer.toString(i), null, - testContextSpecificDocument(), + testContextSpecificDocument("grouping_criteria"), null, // No source, it should be derived null ); @@ -9049,9 +9049,9 @@ public void testShardFailsForCompositeIndexWriterInCaseAddIndexesThrewExceptionW MockDirectoryWrapper wrapper = newMockDirectory(); final Path translogPath = createTempDir("testFailEngineOnRandomIO"); try (Store store = createStore(wrapper)) { - final ParsedDocument doc1 = testParsedDocument("1", null, testContextSpecificDocument(), B_1, null); - final ParsedDocument doc2 = testParsedDocument("2", null, testContextSpecificDocument(), B_1, null); - final ParsedDocument doc3 = testParsedDocument("3", null, testContextSpecificDocument(), B_1, null); + final ParsedDocument doc1 = testParsedDocument("1", null, testContextSpecificDocument("grouping_criteria"), B_1, null); + final ParsedDocument doc2 = testParsedDocument("2", null, testContextSpecificDocument("grouping_criteria"), B_1, null); + final ParsedDocument doc3 = testParsedDocument("3", null, testContextSpecificDocument("grouping_criteria"), B_1, null); AtomicReference throwingIndexWriter = new AtomicReference<>(); final IndexSettings indexSettings = IndexSettingsModule.newIndexSettings( @@ -9094,9 +9094,9 @@ public void testShardFailsForCompositeIndexWriterInCaseAddIndexesThrewExceptionW MockDirectoryWrapper wrapper = newMockDirectory(); final Path translogPath = createTempDir("testFailEngineOnRandomIO"); try (Store store = createStore(wrapper)) { - final ParsedDocument doc1 = testParsedDocument("1", null, testContextSpecificDocument(), B_1, null); - final ParsedDocument doc2 = testParsedDocument("2", null, testContextSpecificDocument(), B_1, null); - final ParsedDocument doc3 = testParsedDocument("1", null, testContextSpecificDocument(), B_1, null); + final ParsedDocument doc1 = testParsedDocument("1", null, testContextSpecificDocument("grouping_criteria"), B_1, null); + final ParsedDocument doc2 = testParsedDocument("2", null, testContextSpecificDocument("grouping_criteria"), B_1, null); + final ParsedDocument doc3 = testParsedDocument("1", null, testContextSpecificDocument("grouping_criteria"), B_1, null); final IndexSettings indexSettings = IndexSettingsModule.newIndexSettings( "test", Settings.builder() @@ -9134,6 +9134,59 @@ public void testShardFailsForCompositeIndexWriterInCaseAddIndexesThrewExceptionW } } + @LockFeatureFlag(CONTEXT_AWARE_MIGRATION_EXPERIMENTAL_FLAG) + public void testDoesNotAllowGroupingCriteriaUpdate() throws IOException, InterruptedException { + final AtomicLong globalCheckpoint = new AtomicLong(SequenceNumbers.NO_OPS_PERFORMED); + final IndexSettings indexSettings = IndexSettingsModule.newIndexSettings( + "test", + Settings.builder() + .put(defaultSettings.getSettings()) + .put(IndexSettings.INDEX_CONTEXT_AWARE_ENABLED_SETTING.getKey(), true) + .build() + ); + try ( + Store store = createStore(); + InternalEngine engine = createEngine( + config(indexSettings, store, createTempDir(), newMergePolicy(), null, null, globalCheckpoint::get) + ) + ) { + final ParsedDocument doc1 = testParsedDocument("1", null, testContextSpecificDocument("grouping_criteria"), B_1, null); + final ParsedDocument doc2 = testParsedDocument("1", null, testContextSpecificDocument("grouping_criteria_update"), B_1, null); + engine.index(indexForDoc(doc1)); + assertThrows(UnsupportedOperationException.class, () -> engine.index(indexForDoc(doc2))); + } + } + + @LockFeatureFlag(CONTEXT_AWARE_MIGRATION_EXPERIMENTAL_FLAG) + public void testAllowGroupingCriteriaUpdateWithTombstone() throws IOException, InterruptedException { + final AtomicLong globalCheckpoint = new AtomicLong(SequenceNumbers.NO_OPS_PERFORMED); + final IndexSettings indexSettings = IndexSettingsModule.newIndexSettings( + "test", + Settings.builder() + .put(defaultSettings.getSettings()) + .put(IndexSettings.INDEX_CONTEXT_AWARE_ENABLED_SETTING.getKey(), true) + .build() + ); + try ( + Store store = createStore(); + InternalEngine engine = createEngine( + config(indexSettings, store, createTempDir(), newMergePolicy(), null, null, globalCheckpoint::get) + ) + ) { + final ParsedDocument doc1 = testParsedDocument("1", null, testContextSpecificDocument("grouping_criteria"), B_1, null); + engine.index(indexForDoc(doc1)); + ParsedDocument doc2 = testParsedDocument("2", null, testContextSpecificDocument("grouping_criteria"), B_1, null); + IndexResult indexResult = engine.index(indexForDoc(doc2)); + doc2 = testParsedDocument("2", null, testContextSpecificDocument("grouping_criteria"), B_1, null); + Engine.Index index = indexForDoc(doc2); + engine.index(index); + engine.delete(new Engine.Delete(index.id(), index.uid(), primaryTerm.get())); + engine.refresh("test"); + ParsedDocument doc3 = testParsedDocument("2", null, testContextSpecificDocument("grouping_criteria_update"), B_1, null); + engine.index(indexForDoc(doc3)); + } + } + private EngineConfig createEngineConfigWithMapperSupplierForDerivedSource(Store store, boolean contextAwareEnabled) throws IOException { // Setup with derived source enabled XContentBuilder mapping; 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 5894ea3abe13f..618fbd0aae59b 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 @@ -379,9 +379,9 @@ protected void assertEngineCleanedUp(Engine engine, TranslogDeletionPolicy trans } } - protected static ParseContext.Document testContextSpecificDocument() { + protected static ParseContext.Document testContextSpecificDocument(String groupingCriteria) { ParseContext.Document doc = testDocumentWithTextField("criteria"); - doc.setGroupingCriteria("grouping_criteria"); + doc.setGroupingCriteria(groupingCriteria); return doc; }