From 7db9af0b8d30ef9b4086994e29387d417f7de992 Mon Sep 17 00:00:00 2001 From: Sagar Upadhyaya Date: Wed, 8 Apr 2026 15:16:15 -0700 Subject: [PATCH 1/9] BitsetFilterCache refactor: Change it to a node level with a size limit Signed-off-by: Sagar Upadhyaya --- .../common/settings/IndexScopedSettings.java | 4 +- .../org/opensearch/index/IndexModule.java | 4 + .../org/opensearch/index/IndexService.java | 16 +- .../opensearch/index/cache/IndexCache.java | 3 +- .../index/cache/bitset/BitsetFilterCache.java | 304 +------------ .../indices/IndicesBitsetFilterCache.java | 413 ++++++++++++++++++ .../opensearch/indices/IndicesService.java | 7 + .../opensearch/index/IndexModuleTests.java | 1 + .../cache/bitset/BitSetFilterCacheTests.java | 79 ++-- .../bucket/nested/NestedAggregatorTests.java | 6 +- .../internal/ContextIndexSearcherTests.java | 7 +- .../search/sort/AbstractSortTestCase.java | 6 +- .../aggregations/AggregatorTestCase.java | 6 +- .../test/AbstractBuilderTestCase.java | 5 +- 14 files changed, 525 insertions(+), 336 deletions(-) create mode 100644 server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java diff --git a/server/src/main/java/org/opensearch/common/settings/IndexScopedSettings.java b/server/src/main/java/org/opensearch/common/settings/IndexScopedSettings.java index 3494f5557d7b3..ac67c1bbb3334 100644 --- a/server/src/main/java/org/opensearch/common/settings/IndexScopedSettings.java +++ b/server/src/main/java/org/opensearch/common/settings/IndexScopedSettings.java @@ -50,7 +50,6 @@ import org.opensearch.index.MergeSchedulerConfig; import org.opensearch.index.SearchSlowLog; import org.opensearch.index.TieredMergePolicyProvider; -import org.opensearch.index.cache.bitset.BitsetFilterCache; import org.opensearch.index.compositeindex.datacube.startree.StarTreeIndexSettings; import org.opensearch.index.engine.EngineConfig; import org.opensearch.index.fielddata.IndexFieldDataService; @@ -59,6 +58,7 @@ import org.opensearch.index.similarity.SimilarityService; import org.opensearch.index.store.FsDirectoryFactory; import org.opensearch.index.store.Store; +import org.opensearch.indices.IndicesBitsetFilterCache; import org.opensearch.indices.IndicesRequestCache; import org.opensearch.search.streaming.FlushModeResolver; @@ -206,7 +206,7 @@ public final class IndexScopedSettings extends AbstractScopedSettings { MapperService.INDEX_MAPPING_TOTAL_FIELDS_LIMIT_SETTING, MapperService.INDEX_MAPPING_DEPTH_LIMIT_SETTING, MapperService.INDEX_MAPPING_FIELD_NAME_LENGTH_LIMIT_SETTING, - BitsetFilterCache.INDEX_LOAD_RANDOM_ACCESS_FILTERS_EAGERLY_SETTING, + IndicesBitsetFilterCache.INDEX_LOAD_RANDOM_ACCESS_FILTERS_EAGERLY_SETTING, IndexModule.INDEX_STORE_TYPE_SETTING, IndexModule.INDEX_COMPOSITE_STORE_TYPE_SETTING, IndexModule.INDEX_STORE_FACTORY_SETTING, diff --git a/server/src/main/java/org/opensearch/index/IndexModule.java b/server/src/main/java/org/opensearch/index/IndexModule.java index c6da5c1b00c8f..4182bf81c535a 100644 --- a/server/src/main/java/org/opensearch/index/IndexModule.java +++ b/server/src/main/java/org/opensearch/index/IndexModule.java @@ -90,6 +90,7 @@ import org.opensearch.index.store.remote.filecache.FileCache; import org.opensearch.index.translog.TranslogFactory; import org.opensearch.indices.ClusterMergeSchedulerConfig; +import org.opensearch.indices.IndicesBitsetFilterCache; import org.opensearch.indices.IndicesQueryCache; import org.opensearch.indices.RemoteStoreSettings; import org.opensearch.indices.fielddata.cache.IndicesFieldDataCache; @@ -763,6 +764,7 @@ public IndexService newIndexService( indicesQueryCache, mapperRegistry, indicesFieldDataCache, + null, namedWriteableRegistry, idFieldDataEnabled, valuesSourceRegistry, @@ -795,6 +797,7 @@ public IndexService newIndexService( IndicesQueryCache indicesQueryCache, MapperRegistry mapperRegistry, IndicesFieldDataCache indicesFieldDataCache, + IndicesBitsetFilterCache indicesBitsetFilterCache, NamedWriteableRegistry namedWriteableRegistry, BooleanSupplier idFieldDataEnabled, ValuesSourceRegistry valuesSourceRegistry, @@ -868,6 +871,7 @@ public IndexService newIndexService( readerWrapperFactory, mapperRegistry, indicesFieldDataCache, + indicesBitsetFilterCache, searchOperationListeners, indexOperationListeners, namedWriteableRegistry, diff --git a/server/src/main/java/org/opensearch/index/IndexService.java b/server/src/main/java/org/opensearch/index/IndexService.java index 79ccb429f9b7b..651986d422805 100644 --- a/server/src/main/java/org/opensearch/index/IndexService.java +++ b/server/src/main/java/org/opensearch/index/IndexService.java @@ -103,6 +103,7 @@ import org.opensearch.index.translog.Translog; import org.opensearch.index.translog.TranslogFactory; import org.opensearch.indices.ClusterMergeSchedulerConfig; +import org.opensearch.indices.IndicesBitsetFilterCache; import org.opensearch.indices.RemoteStoreSettings; import org.opensearch.indices.cluster.IndicesClusterStateService; import org.opensearch.indices.fielddata.cache.IndicesFieldDataCache; @@ -244,6 +245,7 @@ public IndexService( Function> wrapperFactory, MapperRegistry mapperRegistry, IndicesFieldDataCache indicesFieldDataCache, + IndicesBitsetFilterCache indicesBitsetFilterCache, List searchOperationListeners, List indexingOperationListeners, NamedWriteableRegistry namedWriteableRegistry, @@ -315,8 +317,14 @@ public IndexService( this.indexSortSupplier = () -> null; } indexFieldData.setListener(new FieldDataCacheListener(this)); - this.bitsetFilterCache = new BitsetFilterCache(indexSettings, new BitsetCacheListener(this)); - this.warmer = new IndexWarmer(threadPool, indexFieldData, bitsetFilterCache.createListener(threadPool)); + this.bitsetFilterCache = indicesBitsetFilterCache != null + ? new BitsetFilterCache(indicesBitsetFilterCache, new BitsetCacheListener(this)) + : null; + this.warmer = new IndexWarmer( + threadPool, + indexFieldData, + indicesBitsetFilterCache != null ? indicesBitsetFilterCache.createListener(threadPool) : null + ); this.indexCache = new IndexCache(indexSettings, queryCache, bitsetFilterCache); } else { assert indexAnalyzers == null; @@ -450,6 +458,7 @@ public IndexService( wrapperFactory, mapperRegistry, indicesFieldDataCache, + null, searchOperationListeners, indexingOperationListeners, namedWriteableRegistry, @@ -579,7 +588,6 @@ public synchronized void close(final String reason, boolean delete) throws IOExc } } finally { IOUtils.close( - bitsetFilterCache, indexCache, indexFieldData, mapperService, @@ -1138,7 +1146,7 @@ public void accept(ShardLock lock) { * * @opensearch.internal */ - private static final class BitsetCacheListener implements BitsetFilterCache.Listener { + private static final class BitsetCacheListener implements IndicesBitsetFilterCache.Listener { final IndexService indexService; private BitsetCacheListener(IndexService indexService) { diff --git a/server/src/main/java/org/opensearch/index/cache/IndexCache.java b/server/src/main/java/org/opensearch/index/cache/IndexCache.java index 1067863fe9675..04b7641fcf62f 100644 --- a/server/src/main/java/org/opensearch/index/cache/IndexCache.java +++ b/server/src/main/java/org/opensearch/index/cache/IndexCache.java @@ -72,12 +72,11 @@ public BitsetFilterCache bitsetFilterCache() { @Override public void close() throws IOException { - IOUtils.close(queryCache, bitsetFilterCache); + IOUtils.close(queryCache); } public void clear(String reason) { queryCache.clear(reason); - bitsetFilterCache.clear(reason); } } diff --git a/server/src/main/java/org/opensearch/index/cache/bitset/BitsetFilterCache.java b/server/src/main/java/org/opensearch/index/cache/bitset/BitsetFilterCache.java index 96a867c662137..873e39df58a41 100644 --- a/server/src/main/java/org/opensearch/index/cache/bitset/BitsetFilterCache.java +++ b/server/src/main/java/org/opensearch/index/cache/bitset/BitsetFilterCache.java @@ -32,314 +32,36 @@ package org.opensearch.index.cache.bitset; -import org.apache.logging.log4j.message.ParameterizedMessage; -import org.apache.lucene.index.FilterLeafReader; -import org.apache.lucene.index.IndexReader; -import org.apache.lucene.index.IndexReaderContext; -import org.apache.lucene.index.LeafReaderContext; -import org.apache.lucene.index.ReaderUtil; -import org.apache.lucene.search.IndexSearcher; import org.apache.lucene.search.Query; -import org.apache.lucene.search.ScoreMode; -import org.apache.lucene.search.Scorer; -import org.apache.lucene.search.Weight; import org.apache.lucene.search.join.BitSetProducer; -import org.apache.lucene.util.Accountable; -import org.apache.lucene.util.BitDocIdSet; -import org.apache.lucene.util.BitSet; -import org.opensearch.ExceptionsHelper; import org.opensearch.common.annotation.PublicApi; -import org.opensearch.common.cache.Cache; -import org.opensearch.common.cache.CacheBuilder; -import org.opensearch.common.cache.RemovalListener; -import org.opensearch.common.cache.RemovalNotification; -import org.opensearch.common.lucene.index.OpenSearchDirectoryReader; -import org.opensearch.common.lucene.search.Queries; -import org.opensearch.common.settings.Setting; -import org.opensearch.common.settings.Setting.Property; -import org.opensearch.common.unit.TimeValue; -import org.opensearch.core.index.shard.ShardId; -import org.opensearch.index.AbstractIndexComponent; -import org.opensearch.index.IndexSettings; -import org.opensearch.index.IndexWarmer; -import org.opensearch.index.IndexWarmer.TerminationHandle; -import org.opensearch.index.mapper.DocumentMapper; -import org.opensearch.index.mapper.MapperService; -import org.opensearch.index.mapper.ObjectMapper; -import org.opensearch.index.shard.IndexShard; -import org.opensearch.index.shard.ShardUtils; -import org.opensearch.threadpool.ThreadPool; - -import java.io.Closeable; -import java.io.IOException; -import java.util.HashSet; -import java.util.Objects; -import java.util.Set; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.Executor; +import org.opensearch.indices.IndicesBitsetFilterCache; /** - * This is a cache for {@link BitDocIdSet} based filters and is unbounded by size or time. - *

- * Use this cache with care, only components that require that a filter is to be materialized as a {@link BitDocIdSet} - * and require that it should always be around should use this cache, otherwise the - * {@link org.opensearch.index.cache.query.QueryCache} should be used instead. + * Per-index view into the node-level {@link IndicesBitsetFilterCache}. + * Binds a per-index {@link IndicesBitsetFilterCache.Listener} so that callers can obtain + * a {@link BitSetProducer} without needing to supply a listener at every call site. * * @opensearch.api */ @PublicApi(since = "1.0.0") -public final class BitsetFilterCache extends AbstractIndexComponent - implements - IndexReader.ClosedListener, - RemovalListener>, - Closeable { - - public static final Setting INDEX_LOAD_RANDOM_ACCESS_FILTERS_EAGERLY_SETTING = Setting.boolSetting( - "index.load_fixed_bitset_filters_eagerly", - true, - Property.IndexScope - ); +public final class BitsetFilterCache { - private final boolean loadRandomAccessFiltersEagerly; - private final Cache> loadedFilters; - private final Listener listener; + private final IndicesBitsetFilterCache indicesCache; + private final IndicesBitsetFilterCache.Listener listener; - public BitsetFilterCache(IndexSettings indexSettings, Listener listener) { - super(indexSettings); + public BitsetFilterCache(IndicesBitsetFilterCache indicesCache, IndicesBitsetFilterCache.Listener listener) { + if (indicesCache == null) { + throw new IllegalArgumentException("indicesCache must not be null"); + } if (listener == null) { throw new IllegalArgumentException("listener must not be null"); } - this.loadRandomAccessFiltersEagerly = this.indexSettings.getValue(INDEX_LOAD_RANDOM_ACCESS_FILTERS_EAGERLY_SETTING); - this.loadedFilters = CacheBuilder.>builder().removalListener(this).build(); + this.indicesCache = indicesCache; this.listener = listener; } - public static BitSet bitsetFromQuery(Query query, LeafReaderContext context) throws IOException { - final IndexReaderContext topLevelContext = ReaderUtil.getTopLevelContext(context); - final IndexSearcher searcher = new IndexSearcher(topLevelContext); - searcher.setQueryCache(null); - final Weight weight = searcher.createWeight(searcher.rewrite(query), ScoreMode.COMPLETE_NO_SCORES, 1f); - Scorer s = weight.scorer(context); - if (s == null) { - return null; - } else { - return BitSet.of(s.iterator(), context.reader().maxDoc()); - } - } - - public IndexWarmer.Listener createListener(ThreadPool threadPool) { - return new BitSetProducerWarmer(threadPool); - } - public BitSetProducer getBitSetProducer(Query query) { - return new QueryWrapperBitSetProducer(query); - } - - @Override - public void onClose(IndexReader.CacheKey ownerCoreCacheKey) { - loadedFilters.invalidate(ownerCoreCacheKey); - } - - @Override - public void close() { - clear("close"); - } - - public void clear(String reason) { - logger.debug("clearing all bitsets because [{}]", reason); - loadedFilters.invalidateAll(); - } - - private BitSet getAndLoadIfNotPresent(final Query query, final LeafReaderContext context) throws ExecutionException { - final IndexReader.CacheHelper cacheHelper = FilterLeafReader.unwrap(context.reader()).getCoreCacheHelper(); - if (cacheHelper == null) { - throw new IllegalArgumentException("Reader " + context.reader() + " does not support caching"); - } - final IndexReader.CacheKey coreCacheReader = cacheHelper.getKey(); - final ShardId shardId = ShardUtils.extractShardId(context.reader()); - if (indexSettings.getIndex().equals(shardId.getIndex()) == false) { - // insanity - throw new IllegalStateException( - "Trying to load bit set for index " + shardId.getIndex() + " with cache of index " + indexSettings.getIndex() - ); - } - Cache filterToFbs = loadedFilters.computeIfAbsent(coreCacheReader, key -> { - cacheHelper.addClosedListener(BitsetFilterCache.this); - return CacheBuilder.builder().build(); - }); - - return filterToFbs.computeIfAbsent(query, key -> { - final BitSet bitSet = bitsetFromQuery(query, context); - Value value = new Value(bitSet, shardId); - listener.onCache(shardId, value.bitset); - return value; - }).bitset; - } - - @Override - public void onRemoval(RemovalNotification> notification) { - if (notification.getKey() == null) { - return; - } - - Cache valueCache = notification.getValue(); - if (valueCache == null) { - return; - } - - for (Value value : valueCache.values()) { - listener.onRemoval(value.shardId, value.bitset); - // if null then this means the shard has already been removed and the stats are 0 anyway for the shard this key belongs to - } - } - - /** - * Value for bitset filter cache - * - * @opensearch.api - */ - @PublicApi(since = "1.0.0") - public static final class Value { - - final BitSet bitset; - final ShardId shardId; - - public Value(BitSet bitset, ShardId shardId) { - this.bitset = bitset; - this.shardId = shardId; - } - } - - final class QueryWrapperBitSetProducer implements BitSetProducer { - - final Query query; - - QueryWrapperBitSetProducer(Query query) { - this.query = Objects.requireNonNull(query); - } - - // TODO: convertToElastic might need to be renamed - @Override - public BitSet getBitSet(LeafReaderContext context) throws IOException { - try { - return getAndLoadIfNotPresent(query, context); - } catch (ExecutionException e) { - throw ExceptionsHelper.convertToOpenSearchException(e); - } - } - - @Override - public String toString() { - return "random_access(" + query + ")"; - } - - @Override - public boolean equals(Object o) { - if (!(o instanceof QueryWrapperBitSetProducer other)) return false; - return this.query.equals(other.query); - } - - @Override - public int hashCode() { - return 31 * getClass().hashCode() + query.hashCode(); - } - } - - final class BitSetProducerWarmer implements IndexWarmer.Listener { - - private final Executor executor; - - BitSetProducerWarmer(ThreadPool threadPool) { - this.executor = threadPool.executor(ThreadPool.Names.WARMER); - } - - @Override - public IndexWarmer.TerminationHandle warmReader(final IndexShard indexShard, final OpenSearchDirectoryReader reader) { - if (indexSettings.getIndex().equals(indexShard.indexSettings().getIndex()) == false) { - // this is from a different index - return TerminationHandle.NO_WAIT; - } - - if (!loadRandomAccessFiltersEagerly) { - return TerminationHandle.NO_WAIT; - } - - boolean hasNested = false; - final Set warmUp = new HashSet<>(); - final MapperService mapperService = indexShard.mapperService(); - DocumentMapper docMapper = mapperService.documentMapper(); - if (docMapper != null) { - if (docMapper.hasNestedObjects()) { - hasNested = true; - for (ObjectMapper objectMapper : docMapper.objectMappers().values()) { - if (objectMapper.nested().isNested()) { - ObjectMapper parentObjectMapper = objectMapper.getParentObjectMapper(mapperService); - if (parentObjectMapper != null && parentObjectMapper.nested().isNested()) { - warmUp.add(parentObjectMapper.nestedTypeFilter()); - } - } - } - } - } - - if (hasNested) { - warmUp.add(Queries.newNonNestedFilter()); - } - - final CountDownLatch latch = new CountDownLatch(reader.leaves().size() * warmUp.size()); - for (final LeafReaderContext ctx : reader.leaves()) { - for (final Query filterToWarm : warmUp) { - executor.execute(() -> { - try { - final long start = System.nanoTime(); - getAndLoadIfNotPresent(filterToWarm, ctx); - if (indexShard.warmerService().logger().isTraceEnabled()) { - indexShard.warmerService() - .logger() - .trace( - "warmed bitset for [{}], took [{}]", - filterToWarm, - TimeValue.timeValueNanos(System.nanoTime() - start) - ); - } - } catch (Exception e) { - indexShard.warmerService() - .logger() - .warn(() -> new ParameterizedMessage("failed to load " + "bitset for [{}]", filterToWarm), e); - } finally { - latch.countDown(); - } - }); - } - } - return () -> latch.await(); - } - - } - - Cache> getLoadedFilters() { - return loadedFilters; - } - - /** - * A listener interface that is executed for each onCache / onRemoval event - * - * @opensearch.internal - */ - public interface Listener { - /** - * Called for each cached bitset on the cache event. - * @param shardId the shard id the bitset was cached for. This can be null - * @param accountable the bitsets ram representation - */ - void onCache(ShardId shardId, Accountable accountable); - - /** - * Called for each cached bitset on the removal event. - * @param shardId the shard id the bitset was cached for. This can be null - * @param accountable the bitsets ram representation - */ - void onRemoval(ShardId shardId, Accountable accountable); + return indicesCache.getBitSetProducer(query, listener); } } diff --git a/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java b/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java new file mode 100644 index 0000000000000..bfdec7141f181 --- /dev/null +++ b/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java @@ -0,0 +1,413 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +package org.opensearch.indices; + +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; +import org.apache.logging.log4j.message.ParameterizedMessage; +import org.apache.lucene.index.FilterLeafReader; +import org.apache.lucene.index.IndexReader; +import org.apache.lucene.index.IndexReaderContext; +import org.apache.lucene.index.LeafReaderContext; +import org.apache.lucene.index.ReaderUtil; +import org.apache.lucene.search.IndexSearcher; +import org.apache.lucene.search.Query; +import org.apache.lucene.search.ScoreMode; +import org.apache.lucene.search.Scorer; +import org.apache.lucene.search.Weight; +import org.apache.lucene.search.join.BitSetProducer; +import org.apache.lucene.util.Accountable; +import org.apache.lucene.util.BitSet; +import org.opensearch.ExceptionsHelper; +import org.opensearch.common.annotation.PublicApi; +import org.opensearch.common.cache.Cache; +import org.opensearch.common.cache.CacheBuilder; +import org.opensearch.common.cache.RemovalListener; +import org.opensearch.common.cache.RemovalNotification; +import org.opensearch.common.lease.Releasable; +import org.opensearch.common.lucene.index.OpenSearchDirectoryReader; +import org.opensearch.common.lucene.search.Queries; +import org.opensearch.common.settings.Setting; +import org.opensearch.common.settings.Setting.Property; +import org.opensearch.common.settings.Settings; +import org.opensearch.common.unit.TimeValue; +import org.opensearch.common.util.concurrent.ConcurrentCollections; +import org.opensearch.core.common.unit.ByteSizeValue; +import org.opensearch.core.index.shard.ShardId; +import org.opensearch.index.IndexWarmer; +import org.opensearch.index.IndexWarmer.TerminationHandle; +import org.opensearch.index.mapper.DocumentMapper; +import org.opensearch.index.mapper.MapperService; +import org.opensearch.index.mapper.ObjectMapper; +import org.opensearch.index.shard.IndexShard; +import org.opensearch.index.shard.ShardUtils; +import org.opensearch.threadpool.ThreadPool; + +import java.io.Closeable; +import java.io.IOException; +import java.util.HashSet; +import java.util.Objects; +import java.util.Set; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Executor; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.ToLongBiFunction; + +/** + * Node-level cache for {@link BitSet} based filters. Manages a single flat cache shared across + * all indices on the node, with a configurable size limit and async stale entry cleanup. + * + * @opensearch.api + */ +@PublicApi(since = "1.0.0") +public final class IndicesBitsetFilterCache + implements + IndexReader.ClosedListener, + RemovalListener, + Closeable { + + private static final Logger logger = LogManager.getLogger(IndicesBitsetFilterCache.class); + + public static final Setting INDEX_LOAD_RANDOM_ACCESS_FILTERS_EAGERLY_SETTING = Setting.boolSetting( + "index.load_fixed_bitset_filters_eagerly", + true, + Property.IndexScope + ); + + public static final Setting INDICES_BITSET_FILTER_CACHE_SIZE_SETTING = Setting.memorySizeSetting( + "indices.cache.bitset.size", + "5%", + Property.NodeScope + ); + + public static final Setting INDICES_BITSET_FILTER_CACHE_CLEAN_INTERVAL_SETTING = Setting.positiveTimeSetting( + "indices.cache.bitset.cleanup_interval", + TimeValue.timeValueSeconds(60), + Property.NodeScope + ); + + private final Cache cache; + private final Set staleCacheKeys = ConcurrentCollections.newConcurrentSet(); + private final Set registeredKeys = ConcurrentCollections.newConcurrentSet(); + private final BitsetCacheCleaner cacheCleaner; + + public IndicesBitsetFilterCache(Settings settings, ThreadPool threadPool) { + long sizeInBytes = INDICES_BITSET_FILTER_CACHE_SIZE_SETTING.get(settings).getBytes(); + CacheBuilder cacheBuilder = CacheBuilder.builder().removalListener(this); + if (sizeInBytes > 0) { + cacheBuilder.setMaximumWeight(sizeInBytes).weigher(new BitsetWeigher()); + } + this.cache = cacheBuilder.build(); + + TimeValue cleanInterval = INDICES_BITSET_FILTER_CACHE_CLEAN_INTERVAL_SETTING.get(settings); + this.cacheCleaner = new BitsetCacheCleaner(this, threadPool, cleanInterval); + threadPool.schedule(cacheCleaner, cleanInterval, ThreadPool.Names.SAME); + } + + public BitSetProducer getBitSetProducer(Query query, Listener listener) { + return new QueryWrapperBitSetProducer(query, listener); + } + + public IndexWarmer.Listener createListener(ThreadPool threadPool) { + return new BitSetProducerWarmer(threadPool); + } + + public static BitSet bitsetFromQuery(Query query, LeafReaderContext context) throws IOException { + final IndexReaderContext topLevelContext = ReaderUtil.getTopLevelContext(context); + final IndexSearcher searcher = new IndexSearcher(topLevelContext); + searcher.setQueryCache(null); + final Weight weight = searcher.createWeight(searcher.rewrite(query), ScoreMode.COMPLETE_NO_SCORES, 1f); + Scorer s = weight.scorer(context); + if (s == null) { + return null; + } else { + return BitSet.of(s.iterator(), context.reader().maxDoc()); + } + } + + BitSet getAndLoadIfNotPresent(final Query query, final LeafReaderContext context, final Listener listener) throws ExecutionException { + final IndexReader.CacheHelper cacheHelper = FilterLeafReader.unwrap(context.reader()).getCoreCacheHelper(); + if (cacheHelper == null) { + throw new IllegalArgumentException("Reader " + context.reader() + " does not support caching"); + } + final IndexReader.CacheKey coreCacheReader = cacheHelper.getKey(); + final ShardId shardId = ShardUtils.extractShardId(context.reader()); + + if (registeredKeys.add(coreCacheReader)) { + cacheHelper.addClosedListener(this); + } + + final BitsetCacheKey cacheKey = new BitsetCacheKey(coreCacheReader, query); + return cache.computeIfAbsent(cacheKey, key -> { + final BitSet bitSet = bitsetFromQuery(query, context); + Value value = new Value(bitSet, shardId, listener); + listener.onCache(shardId, value.bitset); + return value; + }).bitset; + } + + @Override + public void onClose(IndexReader.CacheKey ownerCoreCacheKey) { + staleCacheKeys.add(ownerCoreCacheKey); + } + + @Override + public void close() { + cacheCleaner.close(); + clear(); + } + + public void clear() { + cache.invalidateAll(); + staleCacheKeys.clear(); + registeredKeys.clear(); + } + + @Override + public void onRemoval(RemovalNotification notification) { + Value value = notification.getValue(); + if (value == null || value.listener == null) { + return; + } + value.listener.onRemoval(value.shardId, value.bitset); + } + + public void purgeStaleEntries() { + if (staleCacheKeys.isEmpty()) { + return; + } + Set staleSnapshot = new HashSet<>(staleCacheKeys); + staleCacheKeys.removeAll(staleSnapshot); + registeredKeys.removeAll(staleSnapshot); + + for (BitsetCacheKey key : cache.keys()) { + if (staleSnapshot.contains(key.readerCacheKey)) { + cache.invalidate(key); + } + } + } + + public Cache getCache() { + return cache; + } + + /** + * Composite key combining a reader segment key with a query. + * + * @opensearch.api + */ + @PublicApi(since = "1.0.0") + public static final class BitsetCacheKey { + final IndexReader.CacheKey readerCacheKey; + final Query query; + + public BitsetCacheKey(IndexReader.CacheKey readerCacheKey, Query query) { + this.readerCacheKey = Objects.requireNonNull(readerCacheKey); + this.query = Objects.requireNonNull(query); + } + + @Override + public boolean equals(Object o) { + if (this == o) return true; + if (!(o instanceof BitsetCacheKey other)) return false; + return readerCacheKey == other.readerCacheKey && query.equals(other.query); + } + + @Override + public int hashCode() { + return 31 * System.identityHashCode(readerCacheKey) + query.hashCode(); + } + } + + /** + * Cached value holding the bitset, shard identity, and the per-index listener for stats. + * + * @opensearch.api + */ + @PublicApi(since = "1.0.0") + public static final class Value { + final BitSet bitset; + final ShardId shardId; + final Listener listener; + + public Value(BitSet bitset, ShardId shardId, Listener listener) { + this.bitset = bitset; + this.shardId = shardId; + this.listener = listener; + } + } + + static class BitsetWeigher implements ToLongBiFunction { + @Override + public long applyAsLong(BitsetCacheKey key, Value value) { + long weight = (value.bitset != null) ? value.bitset.ramBytesUsed() : 0; + return weight == 0 ? 1 : weight; + } + } + + /** + * Listener for per-index cache/removal events, used for shard-level stats tracking. + * + * @opensearch.api + */ + @PublicApi(since = "1.0.0") + public interface Listener { + void onCache(ShardId shardId, Accountable accountable); + + void onRemoval(ShardId shardId, Accountable accountable); + } + + final class QueryWrapperBitSetProducer implements BitSetProducer { + final Query query; + final Listener listener; + + QueryWrapperBitSetProducer(Query query, Listener listener) { + this.query = Objects.requireNonNull(query); + this.listener = Objects.requireNonNull(listener); + } + + @Override + public BitSet getBitSet(LeafReaderContext context) throws IOException { + try { + return getAndLoadIfNotPresent(query, context, listener); + } catch (ExecutionException e) { + throw ExceptionsHelper.convertToOpenSearchException(e); + } + } + + @Override + public String toString() { + return "random_access(" + query + ")"; + } + + @Override + public boolean equals(Object o) { + if (!(o instanceof QueryWrapperBitSetProducer other)) return false; + return this.query.equals(other.query); + } + + @Override + public int hashCode() { + return 31 * getClass().hashCode() + query.hashCode(); + } + } + + final class BitSetProducerWarmer implements IndexWarmer.Listener { + private final Executor executor; + + BitSetProducerWarmer(ThreadPool threadPool) { + this.executor = threadPool.executor(ThreadPool.Names.WARMER); + } + + @Override + public IndexWarmer.TerminationHandle warmReader(final IndexShard indexShard, final OpenSearchDirectoryReader reader) { + if (!indexShard.indexSettings().getValue(INDEX_LOAD_RANDOM_ACCESS_FILTERS_EAGERLY_SETTING)) { + return TerminationHandle.NO_WAIT; + } + + boolean hasNested = false; + final Set warmUp = new HashSet<>(); + final MapperService mapperService = indexShard.mapperService(); + DocumentMapper docMapper = mapperService.documentMapper(); + if (docMapper != null) { + if (docMapper.hasNestedObjects()) { + hasNested = true; + for (ObjectMapper objectMapper : docMapper.objectMappers().values()) { + if (objectMapper.nested().isNested()) { + ObjectMapper parentObjectMapper = objectMapper.getParentObjectMapper(mapperService); + if (parentObjectMapper != null && parentObjectMapper.nested().isNested()) { + warmUp.add(parentObjectMapper.nestedTypeFilter()); + } + } + } + } + } + + if (hasNested) { + warmUp.add(Queries.newNonNestedFilter()); + } + + // Build a listener that routes stats to the correct shard. + final Listener listener = new Listener() { + @Override + public void onCache(ShardId shardId, Accountable accountable) { + if (shardId != null && accountable != null) { + indexShard.shardBitsetFilterCache().onCached(accountable.ramBytesUsed()); + } + } + + @Override + public void onRemoval(ShardId shardId, Accountable accountable) { + if (shardId != null && accountable != null) { + indexShard.shardBitsetFilterCache().onRemoval(accountable.ramBytesUsed()); + } + } + }; + + final CountDownLatch latch = new CountDownLatch(reader.leaves().size() * warmUp.size()); + for (final LeafReaderContext ctx : reader.leaves()) { + for (final Query filterToWarm : warmUp) { + executor.execute(() -> { + try { + final long start = System.nanoTime(); + getAndLoadIfNotPresent(filterToWarm, ctx, listener); + if (indexShard.warmerService().logger().isTraceEnabled()) { + indexShard.warmerService() + .logger() + .trace( + "warmed bitset for [{}], took [{}]", + filterToWarm, + TimeValue.timeValueNanos(System.nanoTime() - start) + ); + } + } catch (Exception e) { + indexShard.warmerService() + .logger() + .warn(() -> new ParameterizedMessage("failed to load bitset for [{}]", filterToWarm), e); + } finally { + latch.countDown(); + } + }); + } + } + return () -> latch.await(); + } + } + + private static final class BitsetCacheCleaner implements Runnable, Releasable { + private final IndicesBitsetFilterCache cache; + private final ThreadPool threadPool; + private final TimeValue interval; + private final AtomicBoolean closed = new AtomicBoolean(false); + + BitsetCacheCleaner(IndicesBitsetFilterCache cache, ThreadPool threadPool, TimeValue interval) { + this.cache = cache; + this.threadPool = threadPool; + this.interval = interval; + } + + @Override + public void run() { + try { + cache.purgeStaleEntries(); + } catch (Exception e) { + logger.warn("Exception during periodic bitset filter cache cleanup:", e); + } + if (closed.get() == false) { + threadPool.scheduleUnlessShuttingDown(interval, ThreadPool.Names.SAME, this); + } + } + + @Override + public void close() { + closed.compareAndSet(false, true); + } + } +} diff --git a/server/src/main/java/org/opensearch/indices/IndicesService.java b/server/src/main/java/org/opensearch/indices/IndicesService.java index ad33e3811f273..56ae5beee8200 100644 --- a/server/src/main/java/org/opensearch/indices/IndicesService.java +++ b/server/src/main/java/org/opensearch/indices/IndicesService.java @@ -377,6 +377,7 @@ public class IndicesService extends AbstractLifecycleComponent private final IndexNameExpressionResolver indexNameExpressionResolver; private final IndexScopedSettings indexScopedSettings; private final IndicesFieldDataCache indicesFieldDataCache; + private final IndicesBitsetFilterCache indicesBitsetFilterCache; private final CacheCleaner cacheCleaner; private final ThreadPool threadPool; private final CircuitBreakerService circuitBreakerService; @@ -523,6 +524,7 @@ public void onRemoval(ShardId shardId, String fieldName, boolean wasEvicted, lon }, clusterService, threadPool); this.cleanInterval = INDICES_CACHE_CLEAN_INTERVAL_SETTING.get(settings); this.cacheCleaner = new CacheCleaner(indicesFieldDataCache, logger, threadPool, this.cleanInterval); + this.indicesBitsetFilterCache = new IndicesBitsetFilterCache(settings, threadPool); this.metaStateService = metaStateService; this.engineFactoryProviders = engineFactoryProviders; @@ -1013,6 +1015,7 @@ public void onStoreClosed(ShardId shardId) { indexMetadata, indicesQueryCache, indicesFieldDataCache, + indicesBitsetFilterCache, finalListeners, indexingMemoryController ); @@ -1068,6 +1071,7 @@ public void onStoreCreated(ShardId shardId) { indexMetadata, indicesQueryCache, indicesFieldDataCache, + indicesBitsetFilterCache, finalListeners, indexingMemoryController ); @@ -1084,6 +1088,7 @@ private synchronized IndexService createIndexService( IndexMetadata indexMetadata, IndicesQueryCache indicesQueryCache, IndicesFieldDataCache indicesFieldDataCache, + IndicesBitsetFilterCache indicesBitsetFilterCache, List builtInListeners, IndexingOperationListener... indexingOperationListeners ) throws IOException { @@ -1139,6 +1144,7 @@ private synchronized IndexService createIndexService( indicesQueryCache, mapperRegistry, indicesFieldDataCache, + indicesBitsetFilterCache, namedWriteableRegistry, this::isIdFieldDataEnabled, valuesSourceRegistry, @@ -1260,6 +1266,7 @@ public synchronized void verifyIndexMetadata(IndexMetadata metadata, IndexMetada metadata, indicesQueryCache, indicesFieldDataCache, + indicesBitsetFilterCache, emptyList() ); closeables.add(() -> service.close("metadata verification", false)); diff --git a/server/src/test/java/org/opensearch/index/IndexModuleTests.java b/server/src/test/java/org/opensearch/index/IndexModuleTests.java index 57ba262b790ea..6db9f152988a2 100644 --- a/server/src/test/java/org/opensearch/index/IndexModuleTests.java +++ b/server/src/test/java/org/opensearch/index/IndexModuleTests.java @@ -268,6 +268,7 @@ private IndexService newIndexService(IndexModule module) throws IOException { indicesQueryCache, mapperRegistry, new IndicesFieldDataCache(settings, listener, clusterService, threadPool), + null, writableRegistry(), () -> false, null, diff --git a/server/src/test/java/org/opensearch/index/cache/bitset/BitSetFilterCacheTests.java b/server/src/test/java/org/opensearch/index/cache/bitset/BitSetFilterCacheTests.java index f3cac6abd6ced..3dda869ac5015 100644 --- a/server/src/test/java/org/opensearch/index/cache/bitset/BitSetFilterCacheTests.java +++ b/server/src/test/java/org/opensearch/index/cache/bitset/BitSetFilterCacheTests.java @@ -54,11 +54,13 @@ import org.opensearch.common.settings.Settings; import org.opensearch.common.util.io.IOUtils; import org.opensearch.core.index.shard.ShardId; -import org.opensearch.index.IndexSettings; -import org.opensearch.test.IndexSettingsModule; +import org.opensearch.indices.IndicesBitsetFilterCache; import org.opensearch.test.OpenSearchTestCase; +import org.opensearch.threadpool.TestThreadPool; +import org.opensearch.threadpool.ThreadPool; import java.io.IOException; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicLong; @@ -66,7 +68,19 @@ public class BitSetFilterCacheTests extends OpenSearchTestCase { - private static final IndexSettings INDEX_SETTINGS = IndexSettingsModule.newIndexSettings("test", Settings.EMPTY); + private ThreadPool threadPool; + + @Override + public void setUp() throws Exception { + super.setUp(); + threadPool = new TestThreadPool("bitset_filter_cache_test"); + } + + @Override + public void tearDown() throws Exception { + ThreadPool.terminate(threadPool, 10, TimeUnit.SECONDS); + super.tearDown(); + } private static int matchCount(BitSetProducer producer, IndexReader reader) throws IOException { int count = 0; @@ -102,24 +116,21 @@ public void testInvalidateEntries() throws Exception { DirectoryReader reader = DirectoryReader.open(writer); reader = OpenSearchDirectoryReader.wrap(reader, new ShardId("test", "_na_", 0)); - BitsetFilterCache cache = new BitsetFilterCache(INDEX_SETTINGS, new BitsetFilterCache.Listener() { + IndicesBitsetFilterCache indicesCache = new IndicesBitsetFilterCache(Settings.EMPTY, threadPool); + BitsetFilterCache cache = new BitsetFilterCache(indicesCache, new IndicesBitsetFilterCache.Listener() { @Override - public void onCache(ShardId shardId, Accountable accountable) { - - } + public void onCache(ShardId shardId, Accountable accountable) {} @Override - public void onRemoval(ShardId shardId, Accountable accountable) { - - } + public void onRemoval(ShardId shardId, Accountable accountable) {} }); BitSetProducer filter = cache.getBitSetProducer(new TermQuery(new Term("field", "value"))); assertThat(matchCount(filter, reader), equalTo(3)); // now cached assertThat(matchCount(filter, reader), equalTo(3)); - // There are 3 segments - assertThat(cache.getLoadedFilters().weight(), equalTo(3L)); + // There are 3 segments, each with 1 query = 3 entries in the flat cache + assertThat(indicesCache.getCache().count(), equalTo(3)); writer.forceMerge(1); reader.close(); @@ -130,13 +141,18 @@ public void onRemoval(ShardId shardId, Accountable accountable) { // now cached assertThat(matchCount(filter, reader), equalTo(3)); - // Only one segment now, so the size must be 1 - assertThat(cache.getLoadedFilters().weight(), equalTo(1L)); + // Old 3 segments were closed (stale entries purged on next access via cleaner), new merged segment cached = 1 + // Trigger purge explicitly since the scheduled cleaner may not have run yet. + indicesCache.purgeStaleEntries(); + assertThat(indicesCache.getCache().count(), equalTo(1)); reader.close(); writer.close(); - // There is no reference from readers and writer to any segment in the test index, so the size in the fbs cache must be 0 - assertThat(cache.getLoadedFilters().weight(), equalTo(0L)); + // Trigger purge for the last closed reader. + indicesCache.purgeStaleEntries(); + assertThat(indicesCache.getCache().count(), equalTo(0)); + + indicesCache.close(); } public void testListener() throws IOException { @@ -155,7 +171,8 @@ public void testListener() throws IOException { final AtomicInteger onCacheCalls = new AtomicInteger(); final AtomicInteger onRemoveCalls = new AtomicInteger(); - final BitsetFilterCache cache = new BitsetFilterCache(INDEX_SETTINGS, new BitsetFilterCache.Listener() { + IndicesBitsetFilterCache indicesCache = new IndicesBitsetFilterCache(Settings.EMPTY, threadPool); + BitsetFilterCache cache = new BitsetFilterCache(indicesCache, new IndicesBitsetFilterCache.Listener() { @Override public void onCache(ShardId shardId, Accountable accountable) { onCacheCalls.incrementAndGet(); @@ -188,31 +205,30 @@ public void onRemoval(ShardId shardId, Accountable accountable) { assertEquals(1, onCacheCalls.get()); assertEquals(0, onRemoveCalls.get()); IOUtils.close(reader, writer); + indicesCache.purgeStaleEntries(); assertEquals(1, onRemoveCalls.get()); assertEquals(0, stats.get()); + + indicesCache.close(); } public void testSetNullListener() { try { - new BitsetFilterCache(INDEX_SETTINGS, null); + new BitsetFilterCache(new IndicesBitsetFilterCache(Settings.EMPTY, threadPool), null); fail("listener can't be null"); } catch (IllegalArgumentException ex) { assertEquals("listener must not be null", ex.getMessage()); - // all is well } } public void testRejectOtherIndex() throws IOException { - BitsetFilterCache cache = new BitsetFilterCache(INDEX_SETTINGS, new BitsetFilterCache.Listener() { + IndicesBitsetFilterCache indicesCache = new IndicesBitsetFilterCache(Settings.EMPTY, threadPool); + BitsetFilterCache cache = new BitsetFilterCache(indicesCache, new IndicesBitsetFilterCache.Listener() { @Override - public void onCache(ShardId shardId, Accountable accountable) { - - } + public void onCache(ShardId shardId, Accountable accountable) {} @Override - public void onRemoval(ShardId shardId, Accountable accountable) { - - } + public void onRemoval(ShardId shardId, Accountable accountable) {} }); Directory dir = newDirectory(); @@ -224,14 +240,17 @@ public void onRemoval(ShardId shardId, Accountable accountable) { BitSetProducer producer = cache.getBitSetProducer(new MatchAllDocsQuery()); + // The node-level cache doesn't validate index identity — it just caches. + // This test verifies the producer works for any index. try { - producer.getBitSet(reader.leaves().get(0)); - fail(); - } catch (IllegalStateException expected) { - assertEquals("Trying to load bit set for index [test2] with cache of index [test]", expected.getMessage()); + for (LeafReaderContext ctx : reader.leaves()) { + producer.getBitSet(ctx); + } } finally { IOUtils.close(reader, dir); } + + indicesCache.close(); } } diff --git a/server/src/test/java/org/opensearch/search/aggregations/bucket/nested/NestedAggregatorTests.java b/server/src/test/java/org/opensearch/search/aggregations/bucket/nested/NestedAggregatorTests.java index c7fbca538c6ee..79801cd953cb0 100644 --- a/server/src/test/java/org/opensearch/search/aggregations/bucket/nested/NestedAggregatorTests.java +++ b/server/src/test/java/org/opensearch/search/aggregations/bucket/nested/NestedAggregatorTests.java @@ -76,6 +76,7 @@ import org.opensearch.index.query.QueryShardContext; import org.opensearch.index.query.TermsQueryBuilder; import org.opensearch.index.query.support.NestedScope; +import org.opensearch.indices.IndicesBitsetFilterCache; import org.opensearch.script.MockScriptEngine; import org.opensearch.script.Script; import org.opensearch.script.ScriptEngine; @@ -1132,7 +1133,10 @@ protected QueryShardContext createQueryShardContext(String fieldName, IndexSetti QueryShardContext queryShardContext = mock(QueryShardContext.class); when(queryShardContext.nestedScope()).thenReturn(new NestedScope(indexSettings)); - BitsetFilterCache bitsetFilterCache = new BitsetFilterCache(indexSettings, Mockito.mock(BitsetFilterCache.Listener.class)); + BitsetFilterCache bitsetFilterCache = new BitsetFilterCache( + Mockito.mock(IndicesBitsetFilterCache.class), + Mockito.mock(IndicesBitsetFilterCache.Listener.class) + ); when(queryShardContext.bitsetFilter(any())).thenReturn(bitsetFilterCache.getBitSetProducer(Queries.newNonNestedFilter())); when(queryShardContext.fieldMapper(anyString())).thenReturn(fieldType); when(queryShardContext.getSearchQuoteAnalyzer(any())).thenCallRealMethod(); diff --git a/server/src/test/java/org/opensearch/search/internal/ContextIndexSearcherTests.java b/server/src/test/java/org/opensearch/search/internal/ContextIndexSearcherTests.java index e806fe764c0a2..460acea0124d0 100644 --- a/server/src/test/java/org/opensearch/search/internal/ContextIndexSearcherTests.java +++ b/server/src/test/java/org/opensearch/search/internal/ContextIndexSearcherTests.java @@ -81,6 +81,7 @@ import org.opensearch.index.cache.bitset.BitsetFilterCache; import org.opensearch.index.shard.IndexShard; import org.opensearch.index.shard.SearchOperationListener; +import org.opensearch.indices.IndicesBitsetFilterCache; import org.opensearch.lucene.util.CombinedBitSet; import org.opensearch.search.aggregations.InternalAggregation; import org.opensearch.search.aggregations.LeafBucketCollector; @@ -90,6 +91,7 @@ import org.opensearch.search.query.QuerySearchResult; import org.opensearch.test.IndexSettingsModule; import org.opensearch.test.OpenSearchTestCase; +import org.opensearch.threadpool.TestThreadPool; import java.io.IOException; import java.io.UncheckedIOException; @@ -244,7 +246,7 @@ public void doTestContextIndexSearcher(boolean sparse, boolean deletions) throws w.deleteDocuments(new Term("delete", "yes")); IndexSettings settings = IndexSettingsModule.newIndexSettings("_index", Settings.EMPTY); - BitsetFilterCache.Listener listener = new BitsetFilterCache.Listener() { + IndicesBitsetFilterCache.Listener listener = new IndicesBitsetFilterCache.Listener() { @Override public void onCache(ShardId shardId, Accountable accountable) { @@ -256,7 +258,8 @@ public void onRemoval(ShardId shardId, Accountable accountable) { } }; DirectoryReader reader = OpenSearchDirectoryReader.wrap(DirectoryReader.open(w), new ShardId(settings.getIndex(), 0)); - BitsetFilterCache cache = new BitsetFilterCache(settings, listener); + IndicesBitsetFilterCache indicesCache = new IndicesBitsetFilterCache(Settings.EMPTY, new TestThreadPool("test")); + BitsetFilterCache cache = new BitsetFilterCache(indicesCache, listener); Query roleQuery = new TermQuery(new Term("allowed", "yes")); BitSet bitSet = cache.getBitSetProducer(roleQuery).getBitSet(reader.leaves().get(0)); if (sparse) { diff --git a/server/src/test/java/org/opensearch/search/sort/AbstractSortTestCase.java b/server/src/test/java/org/opensearch/search/sort/AbstractSortTestCase.java index 257ff1015e3b4..e5bc62bfbc537 100644 --- a/server/src/test/java/org/opensearch/search/sort/AbstractSortTestCase.java +++ b/server/src/test/java/org/opensearch/search/sort/AbstractSortTestCase.java @@ -64,6 +64,7 @@ import org.opensearch.index.query.QueryShardContext; import org.opensearch.index.query.Rewriteable; import org.opensearch.index.query.TermQueryBuilder; +import org.opensearch.indices.IndicesBitsetFilterCache; import org.opensearch.script.MockScriptEngine; import org.opensearch.script.ScriptEngine; import org.opensearch.script.ScriptModule; @@ -208,7 +209,10 @@ protected final QueryShardContext createMockShardContext(IndexSearcher searcher) index, Settings.builder().put(IndexMetadata.SETTING_VERSION_CREATED, Version.CURRENT).build() ); - BitsetFilterCache bitsetFilterCache = new BitsetFilterCache(idxSettings, Mockito.mock(BitsetFilterCache.Listener.class)); + BitsetFilterCache bitsetFilterCache = new BitsetFilterCache( + Mockito.mock(IndicesBitsetFilterCache.class), + Mockito.mock(IndicesBitsetFilterCache.Listener.class) + ); TriFunction, IndexFieldData> indexFieldDataLookup = ( fieldType, fieldIndexName, diff --git a/test/framework/src/main/java/org/opensearch/search/aggregations/AggregatorTestCase.java b/test/framework/src/main/java/org/opensearch/search/aggregations/AggregatorTestCase.java index f2eabf1b8c453..05b3d35a9ec01 100644 --- a/test/framework/src/main/java/org/opensearch/search/aggregations/AggregatorTestCase.java +++ b/test/framework/src/main/java/org/opensearch/search/aggregations/AggregatorTestCase.java @@ -90,7 +90,6 @@ import org.opensearch.index.analysis.IndexAnalyzers; import org.opensearch.index.analysis.NamedAnalyzer; import org.opensearch.index.cache.bitset.BitsetFilterCache; -import org.opensearch.index.cache.bitset.BitsetFilterCache.Listener; import org.opensearch.index.cache.query.DisabledQueryCache; import org.opensearch.index.codec.composite.CompositeIndexFieldInfo; import org.opensearch.index.compositeindex.datacube.Dimension; @@ -128,6 +127,7 @@ import org.opensearch.index.query.QueryShardContext; import org.opensearch.index.shard.IndexShard; import org.opensearch.index.shard.SearchOperationListener; +import org.opensearch.indices.IndicesBitsetFilterCache; import org.opensearch.indices.IndicesModule; import org.opensearch.indices.mapper.MapperRegistry; import org.opensearch.plugins.SearchPlugin; @@ -502,7 +502,9 @@ public boolean shouldCache(Query query) { when(searchContext.numberOfShards()).thenReturn(1); when(searchContext.searcher()).thenReturn(contextIndexSearcher); when(searchContext.fetchPhase()).thenReturn(new FetchPhase(Arrays.asList(new FetchSourcePhase(), new FetchDocValuesPhase()))); - when(searchContext.bitsetFilterCache()).thenReturn(new BitsetFilterCache(indexSettings, mock(Listener.class))); + when(searchContext.bitsetFilterCache()).thenReturn( + new BitsetFilterCache(mock(IndicesBitsetFilterCache.class), mock(IndicesBitsetFilterCache.Listener.class)) + ); IndexShard indexShard = mock(IndexShard.class); when(indexShard.shardId()).thenReturn(new ShardId("test", "test", 0)); when(indexShard.indexSettings()).thenReturn(indexSettings); diff --git a/test/framework/src/main/java/org/opensearch/test/AbstractBuilderTestCase.java b/test/framework/src/main/java/org/opensearch/test/AbstractBuilderTestCase.java index befef4ff57088..4778222823d79 100644 --- a/test/framework/src/main/java/org/opensearch/test/AbstractBuilderTestCase.java +++ b/test/framework/src/main/java/org/opensearch/test/AbstractBuilderTestCase.java @@ -75,6 +75,7 @@ import org.opensearch.index.query.QueryBuilderVisitor; import org.opensearch.index.query.QueryShardContext; import org.opensearch.index.similarity.SimilarityService; +import org.opensearch.indices.IndicesBitsetFilterCache; import org.opensearch.indices.IndicesModule; import org.opensearch.indices.analysis.AnalysisModule; import org.opensearch.indices.fielddata.cache.IndicesFieldDataCache; @@ -380,6 +381,7 @@ private static class ServiceHolder implements Closeable { private final SimilarityService similarityService; private final MapperService mapperService; private final BitsetFilterCache bitsetFilterCache; + private final IndicesBitsetFilterCache indicesBitsetFilterCache; private final ScriptService scriptService; private final Client client; private final long nowInMillis; @@ -450,7 +452,8 @@ private static class ServiceHolder implements Closeable { mapperService, threadPool ); - bitsetFilterCache = new BitsetFilterCache(idxSettings, new BitsetFilterCache.Listener() { + indicesBitsetFilterCache = new IndicesBitsetFilterCache(Settings.EMPTY, threadPool); + bitsetFilterCache = new BitsetFilterCache(indicesBitsetFilterCache, new IndicesBitsetFilterCache.Listener() { @Override public void onCache(ShardId shardId, Accountable accountable) { From 35ab2c01627d91e5b4b7ae0935f492712b205c65 Mon Sep 17 00:00:00 2001 From: Sagar Upadhyaya Date: Wed, 8 Apr 2026 16:21:21 -0700 Subject: [PATCH 2/9] Fix UTs Signed-off-by: Sagar Upadhyaya --- .../indices/IndicesBitsetFilterCache.java | 18 +++++++++--------- .../bucket/nested/NestedAggregatorTests.java | 6 ++++-- .../internal/ContextIndexSearcherTests.java | 7 +++++-- .../search/sort/AbstractSortTestCase.java | 2 +- .../aggregations/AggregatorTestCase.java | 3 ++- .../test/AbstractBuilderTestCase.java | 4 +++- 6 files changed, 24 insertions(+), 16 deletions(-) diff --git a/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java b/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java index bfdec7141f181..a3893cf3928fa 100644 --- a/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java +++ b/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java @@ -25,7 +25,7 @@ import org.apache.lucene.util.Accountable; import org.apache.lucene.util.BitSet; import org.opensearch.ExceptionsHelper; -import org.opensearch.common.annotation.PublicApi; +import org.opensearch.common.annotation.ExperimentalApi; import org.opensearch.common.cache.Cache; import org.opensearch.common.cache.CacheBuilder; import org.opensearch.common.cache.RemovalListener; @@ -66,8 +66,8 @@ * * @opensearch.api */ -@PublicApi(since = "1.0.0") -public final class IndicesBitsetFilterCache +@ExperimentalApi +public class IndicesBitsetFilterCache implements IndexReader.ClosedListener, RemovalListener, @@ -201,9 +201,9 @@ public Cache getCache() { /** * Composite key combining a reader segment key with a query. * - * @opensearch.api + * @opensearch.internal */ - @PublicApi(since = "1.0.0") + @ExperimentalApi public static final class BitsetCacheKey { final IndexReader.CacheKey readerCacheKey; final Query query; @@ -229,9 +229,9 @@ public int hashCode() { /** * Cached value holding the bitset, shard identity, and the per-index listener for stats. * - * @opensearch.api + * @opensearch.internal */ - @PublicApi(since = "1.0.0") + @ExperimentalApi public static final class Value { final BitSet bitset; final ShardId shardId; @@ -255,9 +255,9 @@ public long applyAsLong(BitsetCacheKey key, Value value) { /** * Listener for per-index cache/removal events, used for shard-level stats tracking. * - * @opensearch.api + * @opensearch.internal */ - @PublicApi(since = "1.0.0") + @ExperimentalApi public interface Listener { void onCache(ShardId shardId, Accountable accountable); diff --git a/server/src/test/java/org/opensearch/search/aggregations/bucket/nested/NestedAggregatorTests.java b/server/src/test/java/org/opensearch/search/aggregations/bucket/nested/NestedAggregatorTests.java index 79801cd953cb0..126309b5afdc4 100644 --- a/server/src/test/java/org/opensearch/search/aggregations/bucket/nested/NestedAggregatorTests.java +++ b/server/src/test/java/org/opensearch/search/aggregations/bucket/nested/NestedAggregatorTests.java @@ -49,6 +49,7 @@ import org.apache.lucene.search.DocIdSetIterator; import org.apache.lucene.search.MatchAllDocsQuery; import org.apache.lucene.search.TermQuery; +import org.apache.lucene.search.join.BitSetProducer; import org.apache.lucene.search.join.ScoreMode; import org.apache.lucene.store.Directory; import org.apache.lucene.tests.index.RandomIndexWriter; @@ -1134,10 +1135,11 @@ protected QueryShardContext createQueryShardContext(String fieldName, IndexSetti when(queryShardContext.nestedScope()).thenReturn(new NestedScope(indexSettings)); BitsetFilterCache bitsetFilterCache = new BitsetFilterCache( - Mockito.mock(IndicesBitsetFilterCache.class), + Mockito.mock(IndicesBitsetFilterCache.class, Mockito.RETURNS_DEEP_STUBS), Mockito.mock(IndicesBitsetFilterCache.Listener.class) ); - when(queryShardContext.bitsetFilter(any())).thenReturn(bitsetFilterCache.getBitSetProducer(Queries.newNonNestedFilter())); + BitSetProducer nonNestedFilter = bitsetFilterCache.getBitSetProducer(Queries.newNonNestedFilter()); + when(queryShardContext.bitsetFilter(any())).thenReturn(nonNestedFilter); when(queryShardContext.fieldMapper(anyString())).thenReturn(fieldType); when(queryShardContext.getSearchQuoteAnalyzer(any())).thenCallRealMethod(); when(queryShardContext.getSearchAnalyzer(any())).thenCallRealMethod(); diff --git a/server/src/test/java/org/opensearch/search/internal/ContextIndexSearcherTests.java b/server/src/test/java/org/opensearch/search/internal/ContextIndexSearcherTests.java index 460acea0124d0..6d8dbb5687df3 100644 --- a/server/src/test/java/org/opensearch/search/internal/ContextIndexSearcherTests.java +++ b/server/src/test/java/org/opensearch/search/internal/ContextIndexSearcherTests.java @@ -92,6 +92,7 @@ import org.opensearch.test.IndexSettingsModule; import org.opensearch.test.OpenSearchTestCase; import org.opensearch.threadpool.TestThreadPool; +import org.opensearch.threadpool.ThreadPool; import java.io.IOException; import java.io.UncheckedIOException; @@ -258,7 +259,8 @@ public void onRemoval(ShardId shardId, Accountable accountable) { } }; DirectoryReader reader = OpenSearchDirectoryReader.wrap(DirectoryReader.open(w), new ShardId(settings.getIndex(), 0)); - IndicesBitsetFilterCache indicesCache = new IndicesBitsetFilterCache(Settings.EMPTY, new TestThreadPool("test")); + ThreadPool tp = new TestThreadPool("test"); + IndicesBitsetFilterCache indicesCache = new IndicesBitsetFilterCache(Settings.EMPTY, tp); BitsetFilterCache cache = new BitsetFilterCache(indicesCache, listener); Query roleQuery = new TermQuery(new Term("allowed", "yes")); BitSet bitSet = cache.getBitSetProducer(roleQuery).getBitSet(reader.leaves().get(0)); @@ -318,7 +320,8 @@ public void onRemoval(ShardId shardId, Accountable accountable) { assertEquals(1, topDocs.scoreDocs.length); assertEquals(3f, topDocs.scoreDocs[0].score, 0); - IOUtils.close(reader, w, dir); + IOUtils.close(reader, w, dir, indicesCache); + ThreadPool.terminate(tp, 10, java.util.concurrent.TimeUnit.SECONDS); } public void testSlicesWithMaxTargetSliceSupplier() throws Exception { diff --git a/server/src/test/java/org/opensearch/search/sort/AbstractSortTestCase.java b/server/src/test/java/org/opensearch/search/sort/AbstractSortTestCase.java index e5bc62bfbc537..80cef31c90bc7 100644 --- a/server/src/test/java/org/opensearch/search/sort/AbstractSortTestCase.java +++ b/server/src/test/java/org/opensearch/search/sort/AbstractSortTestCase.java @@ -210,7 +210,7 @@ protected final QueryShardContext createMockShardContext(IndexSearcher searcher) Settings.builder().put(IndexMetadata.SETTING_VERSION_CREATED, Version.CURRENT).build() ); BitsetFilterCache bitsetFilterCache = new BitsetFilterCache( - Mockito.mock(IndicesBitsetFilterCache.class), + Mockito.mock(IndicesBitsetFilterCache.class, Mockito.RETURNS_DEEP_STUBS), Mockito.mock(IndicesBitsetFilterCache.Listener.class) ); TriFunction, IndexFieldData> indexFieldDataLookup = ( diff --git a/test/framework/src/main/java/org/opensearch/search/aggregations/AggregatorTestCase.java b/test/framework/src/main/java/org/opensearch/search/aggregations/AggregatorTestCase.java index 05b3d35a9ec01..38b573ac09761 100644 --- a/test/framework/src/main/java/org/opensearch/search/aggregations/AggregatorTestCase.java +++ b/test/framework/src/main/java/org/opensearch/search/aggregations/AggregatorTestCase.java @@ -181,6 +181,7 @@ import static org.hamcrest.Matchers.instanceOf; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.RETURNS_DEEP_STUBS; import static org.mockito.Mockito.any; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.mock; @@ -503,7 +504,7 @@ public boolean shouldCache(Query query) { when(searchContext.searcher()).thenReturn(contextIndexSearcher); when(searchContext.fetchPhase()).thenReturn(new FetchPhase(Arrays.asList(new FetchSourcePhase(), new FetchDocValuesPhase()))); when(searchContext.bitsetFilterCache()).thenReturn( - new BitsetFilterCache(mock(IndicesBitsetFilterCache.class), mock(IndicesBitsetFilterCache.Listener.class)) + new BitsetFilterCache(mock(IndicesBitsetFilterCache.class, RETURNS_DEEP_STUBS), mock(IndicesBitsetFilterCache.Listener.class)) ); IndexShard indexShard = mock(IndexShard.class); when(indexShard.shardId()).thenReturn(new ShardId("test", "test", 0)); diff --git a/test/framework/src/main/java/org/opensearch/test/AbstractBuilderTestCase.java b/test/framework/src/main/java/org/opensearch/test/AbstractBuilderTestCase.java index 4778222823d79..0e11a932fb5b7 100644 --- a/test/framework/src/main/java/org/opensearch/test/AbstractBuilderTestCase.java +++ b/test/framework/src/main/java/org/opensearch/test/AbstractBuilderTestCase.java @@ -531,7 +531,9 @@ public static Predicate indexNameMatcher() { } @Override - public void close() throws IOException {} + public void close() throws IOException { + indicesBitsetFilterCache.close(); + } QueryShardContext createShardContext(IndexSearcher searcher) { return new QueryShardContext( From 5ceb33c9ad12259aef98ceffc8f39b0529b1d7c0 Mon Sep 17 00:00:00 2001 From: Sagar Upadhyaya Date: Wed, 8 Apr 2026 16:51:08 -0700 Subject: [PATCH 3/9] Fixing build issue after main merge Signed-off-by: Sagar Upadhyaya --- .../org/opensearch/index/IndexModule.java | 2 ++ .../bucket/nested/NestedAggregatorTests.java | 7 ++++++- .../aggregations/AggregatorTestCase.java | 19 +++++++++++++++++-- 3 files changed, 25 insertions(+), 3 deletions(-) diff --git a/server/src/main/java/org/opensearch/index/IndexModule.java b/server/src/main/java/org/opensearch/index/IndexModule.java index 286736da2a26d..cd1046a594968 100644 --- a/server/src/main/java/org/opensearch/index/IndexModule.java +++ b/server/src/main/java/org/opensearch/index/IndexModule.java @@ -839,6 +839,7 @@ public IndexService newIndexService( indicesQueryCache, mapperRegistry, indicesFieldDataCache, + indicesBitsetFilterCache, namedWriteableRegistry, idFieldDataEnabled, valuesSourceRegistry, @@ -871,6 +872,7 @@ public IndexService newIndexService( IndicesQueryCache indicesQueryCache, MapperRegistry mapperRegistry, IndicesFieldDataCache indicesFieldDataCache, + IndicesBitsetFilterCache indicesBitsetFilterCache, NamedWriteableRegistry namedWriteableRegistry, BooleanSupplier idFieldDataEnabled, ValuesSourceRegistry valuesSourceRegistry, diff --git a/server/src/test/java/org/opensearch/search/aggregations/bucket/nested/NestedAggregatorTests.java b/server/src/test/java/org/opensearch/search/aggregations/bucket/nested/NestedAggregatorTests.java index 126309b5afdc4..7eb0eedbf89e2 100644 --- a/server/src/test/java/org/opensearch/search/aggregations/bucket/nested/NestedAggregatorTests.java +++ b/server/src/test/java/org/opensearch/search/aggregations/bucket/nested/NestedAggregatorTests.java @@ -78,6 +78,7 @@ import org.opensearch.index.query.TermsQueryBuilder; import org.opensearch.index.query.support.NestedScope; import org.opensearch.indices.IndicesBitsetFilterCache; +import org.opensearch.threadpool.TestThreadPool; import org.opensearch.script.MockScriptEngine; import org.opensearch.script.Script; import org.opensearch.script.ScriptEngine; @@ -1134,8 +1135,12 @@ protected QueryShardContext createQueryShardContext(String fieldName, IndexSetti QueryShardContext queryShardContext = mock(QueryShardContext.class); when(queryShardContext.nestedScope()).thenReturn(new NestedScope(indexSettings)); + if (aggTestThreadPool == null) { + aggTestThreadPool = new TestThreadPool("nested_agg_test"); + aggTestIndicesBitsetFilterCache = new IndicesBitsetFilterCache(Settings.EMPTY, aggTestThreadPool); + } BitsetFilterCache bitsetFilterCache = new BitsetFilterCache( - Mockito.mock(IndicesBitsetFilterCache.class, Mockito.RETURNS_DEEP_STUBS), + aggTestIndicesBitsetFilterCache, Mockito.mock(IndicesBitsetFilterCache.Listener.class) ); BitSetProducer nonNestedFilter = bitsetFilterCache.getBitSetProducer(Queries.newNonNestedFilter()); diff --git a/test/framework/src/main/java/org/opensearch/search/aggregations/AggregatorTestCase.java b/test/framework/src/main/java/org/opensearch/search/aggregations/AggregatorTestCase.java index 38b573ac09761..19d7dfa876f27 100644 --- a/test/framework/src/main/java/org/opensearch/search/aggregations/AggregatorTestCase.java +++ b/test/framework/src/main/java/org/opensearch/search/aggregations/AggregatorTestCase.java @@ -128,6 +128,8 @@ import org.opensearch.index.shard.IndexShard; import org.opensearch.index.shard.SearchOperationListener; import org.opensearch.indices.IndicesBitsetFilterCache; +import org.opensearch.threadpool.TestThreadPool; +import org.opensearch.threadpool.ThreadPool; import org.opensearch.indices.IndicesModule; import org.opensearch.indices.mapper.MapperRegistry; import org.opensearch.plugins.SearchPlugin; @@ -181,7 +183,6 @@ import static org.hamcrest.Matchers.instanceOf; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyString; -import static org.mockito.Mockito.RETURNS_DEEP_STUBS; import static org.mockito.Mockito.any; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.mock; @@ -197,6 +198,8 @@ public abstract class AggregatorTestCase extends OpenSearchTestCase { private List releasables = new ArrayList<>(); private static final String TYPE_NAME = "type"; protected ValuesSourceRegistry valuesSourceRegistry; + protected ThreadPool aggTestThreadPool; + protected IndicesBitsetFilterCache aggTestIndicesBitsetFilterCache; // A list of field types that should not be tested, or are not currently supported private static List TYPE_TEST_DENYLIST; @@ -503,8 +506,12 @@ public boolean shouldCache(Query query) { when(searchContext.numberOfShards()).thenReturn(1); when(searchContext.searcher()).thenReturn(contextIndexSearcher); when(searchContext.fetchPhase()).thenReturn(new FetchPhase(Arrays.asList(new FetchSourcePhase(), new FetchDocValuesPhase()))); + if (aggTestThreadPool == null) { + aggTestThreadPool = new TestThreadPool("agg_test"); + aggTestIndicesBitsetFilterCache = new IndicesBitsetFilterCache(Settings.EMPTY, aggTestThreadPool); + } when(searchContext.bitsetFilterCache()).thenReturn( - new BitsetFilterCache(mock(IndicesBitsetFilterCache.class, RETURNS_DEEP_STUBS), mock(IndicesBitsetFilterCache.Listener.class)) + new BitsetFilterCache(aggTestIndicesBitsetFilterCache, mock(IndicesBitsetFilterCache.Listener.class)) ); IndexShard indexShard = mock(IndexShard.class); when(indexShard.shardId()).thenReturn(new ShardId("test", "test", 0)); @@ -1381,6 +1388,14 @@ public IndexAnalyzers getIndexAnalyzers() { private void cleanupReleasables() { Releasables.close(releasables); releasables.clear(); + if (aggTestIndicesBitsetFilterCache != null) { + aggTestIndicesBitsetFilterCache.close(); + aggTestIndicesBitsetFilterCache = null; + } + if (aggTestThreadPool != null) { + ThreadPool.terminate(aggTestThreadPool, 10, java.util.concurrent.TimeUnit.SECONDS); + aggTestThreadPool = null; + } } /** From 21f368ac8690ea8ffcbb682b8efeabb55d5bf9fc Mon Sep 17 00:00:00 2001 From: Sagar Upadhyaya Date: Thu, 9 Apr 2026 15:57:29 -0700 Subject: [PATCH 4/9] Fix backward compatibility issue Signed-off-by: Sagar Upadhyaya --- .../org/opensearch/index/IndexModule.java | 72 ++++ .../org/opensearch/index/IndexService.java | 5 +- .../opensearch/index/cache/IndexCache.java | 3 +- .../index/cache/bitset/BitsetFilterCache.java | 128 ++++++- .../indices/IndicesBitsetFilterCache.java | 33 +- .../cache/bitset/BitSetFilterCacheTests.java | 11 +- .../IndicesBitsetFilterCacheTests.java | 322 ++++++++++++++++++ .../bucket/nested/NestedAggregatorTests.java | 5 +- .../internal/ContextIndexSearcherTests.java | 4 +- .../search/sort/AbstractSortTestCase.java | 3 +- .../aggregations/AggregatorTestCase.java | 6 +- .../test/AbstractBuilderTestCase.java | 2 +- 12 files changed, 545 insertions(+), 49 deletions(-) create mode 100644 server/src/test/java/org/opensearch/indices/IndicesBitsetFilterCacheTests.java diff --git a/server/src/main/java/org/opensearch/index/IndexModule.java b/server/src/main/java/org/opensearch/index/IndexModule.java index cd1046a594968..e30983de5df7b 100644 --- a/server/src/main/java/org/opensearch/index/IndexModule.java +++ b/server/src/main/java/org/opensearch/index/IndexModule.java @@ -784,6 +784,78 @@ public IndexService newIndexService( ); } + /** + * @deprecated Use the overload that accepts {@code indicesBitsetFilterCache} and {@code dataFormatRegistry} parameters. + */ + @Deprecated(forRemoval = true) + public IndexService newIndexService( + IndexService.IndexCreationContext indexCreationContext, + NodeEnvironment environment, + NamedXContentRegistry xContentRegistry, + IndexService.ShardStoreDeleter shardStoreDeleter, + CircuitBreakerService circuitBreakerService, + BigArrays bigArrays, + ThreadPool threadPool, + ScriptService scriptService, + ClusterService clusterService, + Client client, + IndicesQueryCache indicesQueryCache, + MapperRegistry mapperRegistry, + IndicesFieldDataCache indicesFieldDataCache, + NamedWriteableRegistry namedWriteableRegistry, + BooleanSupplier idFieldDataEnabled, + ValuesSourceRegistry valuesSourceRegistry, + IndexStorePlugin.DirectoryFactory remoteDirectoryFactory, + BiFunction translogFactorySupplier, + Supplier clusterDefaultRefreshIntervalSupplier, + Supplier fixedRefreshIntervalSchedulingEnabled, + Supplier shardLevelRefreshEnabled, + RecoverySettings recoverySettings, + RemoteStoreSettings remoteStoreSettings, + Consumer replicator, + Function segmentReplicationStatsProvider, + Supplier clusterDefaultMaxMergeAtOnceSupplier, + ClusterMergeSchedulerConfig clusterMergeSchedulerConfig, + CheckedTriFunction< + ShardPath, + MapperService, + IndexSettings, + DataFormatAwareEngineFactory, + IOException> dataFormatAwareEngineFactorySupplier + ) throws IOException { + return newIndexService( + indexCreationContext, + environment, + xContentRegistry, + shardStoreDeleter, + circuitBreakerService, + bigArrays, + threadPool, + scriptService, + clusterService, + client, + indicesQueryCache, + mapperRegistry, + indicesFieldDataCache, + null, + namedWriteableRegistry, + idFieldDataEnabled, + valuesSourceRegistry, + remoteDirectoryFactory, + translogFactorySupplier, + clusterDefaultRefreshIntervalSupplier, + fixedRefreshIntervalSchedulingEnabled, + shardLevelRefreshEnabled, + recoverySettings, + remoteStoreSettings, + replicator, + segmentReplicationStatsProvider, + clusterDefaultMaxMergeAtOnceSupplier, + clusterMergeSchedulerConfig, + (DataFormatRegistry) null + ); + } + /** * @deprecated Use the overload that accepts a {@code dataFormatRegistry} parameter. */ diff --git a/server/src/main/java/org/opensearch/index/IndexService.java b/server/src/main/java/org/opensearch/index/IndexService.java index 1055448092570..9f2d4bb157ac3 100644 --- a/server/src/main/java/org/opensearch/index/IndexService.java +++ b/server/src/main/java/org/opensearch/index/IndexService.java @@ -307,7 +307,7 @@ public IndexService( } indexFieldData.setListener(new FieldDataCacheListener(this)); this.bitsetFilterCache = indicesBitsetFilterCache != null - ? new BitsetFilterCache(indicesBitsetFilterCache, new BitsetCacheListener(this)) + ? new BitsetFilterCache(indexSettings, indicesBitsetFilterCache, new BitsetCacheListener(this)) : null; this.warmer = new IndexWarmer( threadPool, @@ -577,6 +577,7 @@ public synchronized void close(final String reason, boolean delete) throws IOExc } } finally { IOUtils.close( + bitsetFilterCache, indexCache, indexFieldData, mapperService, @@ -1132,7 +1133,7 @@ public void accept(ShardLock lock) { * * @opensearch.internal */ - private static final class BitsetCacheListener implements IndicesBitsetFilterCache.Listener { + private static final class BitsetCacheListener implements BitsetFilterCache.Listener { final IndexService indexService; private BitsetCacheListener(IndexService indexService) { diff --git a/server/src/main/java/org/opensearch/index/cache/IndexCache.java b/server/src/main/java/org/opensearch/index/cache/IndexCache.java index 04b7641fcf62f..1067863fe9675 100644 --- a/server/src/main/java/org/opensearch/index/cache/IndexCache.java +++ b/server/src/main/java/org/opensearch/index/cache/IndexCache.java @@ -72,11 +72,12 @@ public BitsetFilterCache bitsetFilterCache() { @Override public void close() throws IOException { - IOUtils.close(queryCache); + IOUtils.close(queryCache, bitsetFilterCache); } public void clear(String reason) { queryCache.clear(reason); + bitsetFilterCache.clear(reason); } } diff --git a/server/src/main/java/org/opensearch/index/cache/bitset/BitsetFilterCache.java b/server/src/main/java/org/opensearch/index/cache/bitset/BitsetFilterCache.java index 873e39df58a41..e9527ad666c93 100644 --- a/server/src/main/java/org/opensearch/index/cache/bitset/BitsetFilterCache.java +++ b/server/src/main/java/org/opensearch/index/cache/bitset/BitsetFilterCache.java @@ -32,28 +32,61 @@ package org.opensearch.index.cache.bitset; +import org.apache.lucene.index.IndexReader; +import org.apache.lucene.index.LeafReaderContext; import org.apache.lucene.search.Query; import org.apache.lucene.search.join.BitSetProducer; +import org.apache.lucene.util.Accountable; +import org.apache.lucene.util.BitDocIdSet; +import org.apache.lucene.util.BitSet; import org.opensearch.common.annotation.PublicApi; +import org.opensearch.common.cache.Cache; +import org.opensearch.common.cache.RemovalListener; +import org.opensearch.common.cache.RemovalNotification; +import org.opensearch.common.settings.Setting; +import org.opensearch.core.index.shard.ShardId; +import org.opensearch.index.AbstractIndexComponent; +import org.opensearch.index.IndexSettings; +import org.opensearch.index.IndexWarmer; import org.opensearch.indices.IndicesBitsetFilterCache; +import org.opensearch.threadpool.ThreadPool; + +import java.io.Closeable; +import java.io.IOException; /** * Per-index view into the node-level {@link IndicesBitsetFilterCache}. - * Binds a per-index {@link IndicesBitsetFilterCache.Listener} so that callers can obtain - * a {@link BitSetProducer} without needing to supply a listener at every call site. + *

+ * This is a cache for {@link BitDocIdSet} based filters. Use this cache with care, only components + * that require a filter to be materialized as a {@link BitDocIdSet} and require that it should always + * be around should use this cache, otherwise the + * {@link org.opensearch.index.cache.query.QueryCache} should be used instead. * * @opensearch.api */ @PublicApi(since = "1.0.0") -public final class BitsetFilterCache { +public final class BitsetFilterCache extends AbstractIndexComponent + implements + IndexReader.ClosedListener, + RemovalListener>, + Closeable { + + public static final Setting INDEX_LOAD_RANDOM_ACCESS_FILTERS_EAGERLY_SETTING = + IndicesBitsetFilterCache.INDEX_LOAD_RANDOM_ACCESS_FILTERS_EAGERLY_SETTING; private final IndicesBitsetFilterCache indicesCache; - private final IndicesBitsetFilterCache.Listener listener; + private final Listener listener; - public BitsetFilterCache(IndicesBitsetFilterCache indicesCache, IndicesBitsetFilterCache.Listener listener) { - if (indicesCache == null) { - throw new IllegalArgumentException("indicesCache must not be null"); - } + /** + * @deprecated Use {@link #BitsetFilterCache(IndexSettings, IndicesBitsetFilterCache, Listener)} instead. + */ + @Deprecated + public BitsetFilterCache(IndexSettings indexSettings, Listener listener) { + this(indexSettings, null, listener); + } + + public BitsetFilterCache(IndexSettings indexSettings, IndicesBitsetFilterCache indicesCache, Listener listener) { + super(indexSettings); if (listener == null) { throw new IllegalArgumentException("listener must not be null"); } @@ -61,7 +94,84 @@ public BitsetFilterCache(IndicesBitsetFilterCache indicesCache, IndicesBitsetFil this.listener = listener; } + public static BitSet bitsetFromQuery(Query query, LeafReaderContext context) throws IOException { + return IndicesBitsetFilterCache.bitsetFromQuery(query, context); + } + + /** + * @deprecated The warmer is now created by {@link IndicesBitsetFilterCache#createListener(ThreadPool)}. + */ + @Deprecated + public IndexWarmer.Listener createListener(ThreadPool threadPool) { + if (indicesCache != null) { + return indicesCache.createListener(threadPool); + } + return null; + } + public BitSetProducer getBitSetProducer(Query query) { - return indicesCache.getBitSetProducer(query, listener); + if (indicesCache != null) { + return indicesCache.getBitSetProducer(query, listener); + } + throw new IllegalStateException("IndicesBitsetFilterCache is not available"); + } + + @Override + public void onClose(IndexReader.CacheKey ownerCoreCacheKey) { + // Delegated to node-level cache + } + + @Override + public void close() { + // Per-index close is a no-op; the node-level cache manages the lifecycle. + } + + public void clear(String reason) { + logger.debug("clearing all bitsets because [{}]", reason); + // Per-index clear is a no-op; entries are evicted by the node-level cache. + } + + @Override + public void onRemoval(RemovalNotification> notification) { + // Delegated to node-level cache + } + + /** + * Value for bitset filter cache + * + * @opensearch.api + */ + @PublicApi(since = "1.0.0") + public static final class Value { + + final BitSet bitset; + final ShardId shardId; + + public Value(BitSet bitset, ShardId shardId) { + this.bitset = bitset; + this.shardId = shardId; + } + } + + /** + * A listener interface that is executed for each onCache / onRemoval event + * + * @opensearch.internal + */ + @PublicApi(since = "1.0.0") + public interface Listener { + /** + * Called for each cached bitset on the cache event. + * @param shardId the shard id the bitset was cached for. This can be null + * @param accountable the bitsets ram representation + */ + void onCache(ShardId shardId, Accountable accountable); + + /** + * Called for each cached bitset on the removal event. + * @param shardId the shard id the bitset was cached for. This can be null + * @param accountable the bitsets ram representation + */ + void onRemoval(ShardId shardId, Accountable accountable); } } diff --git a/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java b/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java index a3893cf3928fa..99e72bc8826e4 100644 --- a/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java +++ b/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java @@ -42,6 +42,7 @@ import org.opensearch.core.index.shard.ShardId; import org.opensearch.index.IndexWarmer; import org.opensearch.index.IndexWarmer.TerminationHandle; +import org.opensearch.index.cache.bitset.BitsetFilterCache; import org.opensearch.index.mapper.DocumentMapper; import org.opensearch.index.mapper.MapperService; import org.opensearch.index.mapper.ObjectMapper; @@ -111,7 +112,7 @@ public IndicesBitsetFilterCache(Settings settings, ThreadPool threadPool) { threadPool.schedule(cacheCleaner, cleanInterval, ThreadPool.Names.SAME); } - public BitSetProducer getBitSetProducer(Query query, Listener listener) { + public BitSetProducer getBitSetProducer(Query query, BitsetFilterCache.Listener listener) { return new QueryWrapperBitSetProducer(query, listener); } @@ -132,7 +133,8 @@ public static BitSet bitsetFromQuery(Query query, LeafReaderContext context) thr } } - BitSet getAndLoadIfNotPresent(final Query query, final LeafReaderContext context, final Listener listener) throws ExecutionException { + BitSet getAndLoadIfNotPresent(final Query query, final LeafReaderContext context, final BitsetFilterCache.Listener listener) + throws ExecutionException { final IndexReader.CacheHelper cacheHelper = FilterLeafReader.unwrap(context.reader()).getCoreCacheHelper(); if (cacheHelper == null) { throw new IllegalArgumentException("Reader " + context.reader() + " does not support caching"); @@ -226,18 +228,13 @@ public int hashCode() { } } - /** - * Cached value holding the bitset, shard identity, and the per-index listener for stats. - * - * @opensearch.internal - */ @ExperimentalApi public static final class Value { final BitSet bitset; final ShardId shardId; - final Listener listener; + final BitsetFilterCache.Listener listener; - public Value(BitSet bitset, ShardId shardId, Listener listener) { + Value(BitSet bitset, ShardId shardId, BitsetFilterCache.Listener listener) { this.bitset = bitset; this.shardId = shardId; this.listener = listener; @@ -252,23 +249,11 @@ public long applyAsLong(BitsetCacheKey key, Value value) { } } - /** - * Listener for per-index cache/removal events, used for shard-level stats tracking. - * - * @opensearch.internal - */ - @ExperimentalApi - public interface Listener { - void onCache(ShardId shardId, Accountable accountable); - - void onRemoval(ShardId shardId, Accountable accountable); - } - final class QueryWrapperBitSetProducer implements BitSetProducer { final Query query; - final Listener listener; + final BitsetFilterCache.Listener listener; - QueryWrapperBitSetProducer(Query query, Listener listener) { + QueryWrapperBitSetProducer(Query query, BitsetFilterCache.Listener listener) { this.query = Objects.requireNonNull(query); this.listener = Objects.requireNonNull(listener); } @@ -335,7 +320,7 @@ public IndexWarmer.TerminationHandle warmReader(final IndexShard indexShard, fin } // Build a listener that routes stats to the correct shard. - final Listener listener = new Listener() { + final BitsetFilterCache.Listener listener = new BitsetFilterCache.Listener() { @Override public void onCache(ShardId shardId, Accountable accountable) { if (shardId != null && accountable != null) { diff --git a/server/src/test/java/org/opensearch/index/cache/bitset/BitSetFilterCacheTests.java b/server/src/test/java/org/opensearch/index/cache/bitset/BitSetFilterCacheTests.java index 3dda869ac5015..f65a2d18b5970 100644 --- a/server/src/test/java/org/opensearch/index/cache/bitset/BitSetFilterCacheTests.java +++ b/server/src/test/java/org/opensearch/index/cache/bitset/BitSetFilterCacheTests.java @@ -54,7 +54,9 @@ import org.opensearch.common.settings.Settings; import org.opensearch.common.util.io.IOUtils; import org.opensearch.core.index.shard.ShardId; +import org.opensearch.index.IndexSettings; import org.opensearch.indices.IndicesBitsetFilterCache; +import org.opensearch.test.IndexSettingsModule; import org.opensearch.test.OpenSearchTestCase; import org.opensearch.threadpool.TestThreadPool; import org.opensearch.threadpool.ThreadPool; @@ -68,6 +70,7 @@ public class BitSetFilterCacheTests extends OpenSearchTestCase { + private static final IndexSettings INDEX_SETTINGS = IndexSettingsModule.newIndexSettings("test", Settings.EMPTY); private ThreadPool threadPool; @Override @@ -117,7 +120,7 @@ public void testInvalidateEntries() throws Exception { reader = OpenSearchDirectoryReader.wrap(reader, new ShardId("test", "_na_", 0)); IndicesBitsetFilterCache indicesCache = new IndicesBitsetFilterCache(Settings.EMPTY, threadPool); - BitsetFilterCache cache = new BitsetFilterCache(indicesCache, new IndicesBitsetFilterCache.Listener() { + BitsetFilterCache cache = new BitsetFilterCache(INDEX_SETTINGS, indicesCache, new BitsetFilterCache.Listener() { @Override public void onCache(ShardId shardId, Accountable accountable) {} @@ -172,7 +175,7 @@ public void testListener() throws IOException { final AtomicInteger onRemoveCalls = new AtomicInteger(); IndicesBitsetFilterCache indicesCache = new IndicesBitsetFilterCache(Settings.EMPTY, threadPool); - BitsetFilterCache cache = new BitsetFilterCache(indicesCache, new IndicesBitsetFilterCache.Listener() { + BitsetFilterCache cache = new BitsetFilterCache(INDEX_SETTINGS, indicesCache, new BitsetFilterCache.Listener() { @Override public void onCache(ShardId shardId, Accountable accountable) { onCacheCalls.incrementAndGet(); @@ -214,7 +217,7 @@ public void onRemoval(ShardId shardId, Accountable accountable) { public void testSetNullListener() { try { - new BitsetFilterCache(new IndicesBitsetFilterCache(Settings.EMPTY, threadPool), null); + new BitsetFilterCache(INDEX_SETTINGS, new IndicesBitsetFilterCache(Settings.EMPTY, threadPool), null); fail("listener can't be null"); } catch (IllegalArgumentException ex) { assertEquals("listener must not be null", ex.getMessage()); @@ -223,7 +226,7 @@ public void testSetNullListener() { public void testRejectOtherIndex() throws IOException { IndicesBitsetFilterCache indicesCache = new IndicesBitsetFilterCache(Settings.EMPTY, threadPool); - BitsetFilterCache cache = new BitsetFilterCache(indicesCache, new IndicesBitsetFilterCache.Listener() { + BitsetFilterCache cache = new BitsetFilterCache(INDEX_SETTINGS, indicesCache, new BitsetFilterCache.Listener() { @Override public void onCache(ShardId shardId, Accountable accountable) {} diff --git a/server/src/test/java/org/opensearch/indices/IndicesBitsetFilterCacheTests.java b/server/src/test/java/org/opensearch/indices/IndicesBitsetFilterCacheTests.java new file mode 100644 index 0000000000000..a1a6891589973 --- /dev/null +++ b/server/src/test/java/org/opensearch/indices/IndicesBitsetFilterCacheTests.java @@ -0,0 +1,322 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +package org.opensearch.indices; + +import org.apache.lucene.analysis.standard.StandardAnalyzer; +import org.apache.lucene.document.Document; +import org.apache.lucene.document.Field; +import org.apache.lucene.document.StringField; +import org.apache.lucene.index.DirectoryReader; +import org.apache.lucene.index.IndexWriter; +import org.apache.lucene.index.IndexWriterConfig; +import org.apache.lucene.index.LeafReaderContext; +import org.apache.lucene.index.LogByteSizeMergePolicy; +import org.apache.lucene.index.Term; +import org.apache.lucene.search.TermQuery; +import org.apache.lucene.search.join.BitSetProducer; +import org.apache.lucene.store.ByteBuffersDirectory; +import org.apache.lucene.util.Accountable; +import org.opensearch.common.lucene.index.OpenSearchDirectoryReader; +import org.opensearch.common.settings.Settings; +import org.opensearch.core.index.shard.ShardId; +import org.opensearch.index.cache.bitset.BitsetFilterCache; +import org.opensearch.test.OpenSearchTestCase; +import org.opensearch.threadpool.TestThreadPool; +import org.opensearch.threadpool.ThreadPool; + +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; + +import static org.hamcrest.Matchers.equalTo; +import static org.hamcrest.Matchers.greaterThan; +import static org.hamcrest.Matchers.lessThanOrEqualTo; + +public class IndicesBitsetFilterCacheTests extends OpenSearchTestCase { + + private ThreadPool threadPool; + + @Override + public void setUp() throws Exception { + super.setUp(); + threadPool = new TestThreadPool("indices_bitset_filter_cache_test"); + } + + @Override + public void tearDown() throws Exception { + ThreadPool.terminate(threadPool, 10, TimeUnit.SECONDS); + super.tearDown(); + } + + private static final BitsetFilterCache.Listener NO_OP_LISTENER = new BitsetFilterCache.Listener() { + @Override + public void onCache(ShardId shardId, Accountable accountable) {} + + @Override + public void onRemoval(ShardId shardId, Accountable accountable) {} + }; + + /** + * Verifies that when the cache size setting is configured, the cache evicts entries + * once the total weight exceeds the configured limit. + */ + public void testCacheSizeLimitIsHonored() throws Exception { + // First, figure out how large a single bitset entry is by caching one. + long singleEntryBytes; + try (IndicesBitsetFilterCache probeCache = new IndicesBitsetFilterCache(Settings.EMPTY, threadPool)) { + IndexWriter writer = new IndexWriter( + new ByteBuffersDirectory(), + new IndexWriterConfig(new StandardAnalyzer()).setMergePolicy(new LogByteSizeMergePolicy()) + ); + Document doc = new Document(); + doc.add(new StringField("field", "value", Field.Store.NO)); + writer.addDocument(doc); + writer.commit(); + + DirectoryReader reader = DirectoryReader.open(writer); + reader = OpenSearchDirectoryReader.wrap(reader, new ShardId("probe", "_na_", 0)); + + BitSetProducer producer = probeCache.getBitSetProducer(new TermQuery(new Term("field", "value")), NO_OP_LISTENER); + producer.getBitSet(reader.leaves().get(0)); + + assertThat(probeCache.getCache().count(), equalTo(1)); + singleEntryBytes = probeCache.getCache().weight(); + assertThat(singleEntryBytes, greaterThan(0L)); + + reader.close(); + writer.close(); + } + + // Now create a cache with a size limit that fits exactly 2 entries. + long cacheSizeBytes = singleEntryBytes * 2; + Settings settings = Settings.builder().put("indices.cache.bitset.size", cacheSizeBytes + "b").build(); + + try (IndicesBitsetFilterCache cache = new IndicesBitsetFilterCache(settings, threadPool)) { + // Create an index with 3 separate segments (3 commits), each with one doc matching a different query. + IndexWriter writer = new IndexWriter( + new ByteBuffersDirectory(), + new IndexWriterConfig(new StandardAnalyzer()).setMergePolicy(new LogByteSizeMergePolicy()) + ); + + for (int i = 0; i < 3; i++) { + Document doc = new Document(); + doc.add(new StringField("field", "val" + i, Field.Store.NO)); + writer.addDocument(doc); + writer.commit(); + } + + DirectoryReader reader = DirectoryReader.open(writer); + reader = OpenSearchDirectoryReader.wrap(reader, new ShardId("test", "_na_", 0)); + + // 3 segments, cache 3 different queries — one per segment. + // Each segment has exactly 1 leaf. + assertThat(reader.leaves().size(), equalTo(3)); + + for (int i = 0; i < 3; i++) { + LeafReaderContext leaf = reader.leaves().get(i); + BitSetProducer producer = cache.getBitSetProducer(new TermQuery(new Term("field", "val" + i)), NO_OP_LISTENER); + producer.getBitSet(leaf); + } + + // We inserted 3 entries but the cache can only hold 2. + // The LRU eviction should have kicked in. + assertThat(cache.getCache().count(), lessThanOrEqualTo(2)); + assertThat(cache.getCache().weight(), lessThanOrEqualTo(cacheSizeBytes)); + + reader.close(); + writer.close(); + } + } + + /** + * Verifies that the onRemoval listener is called when entries are evicted due to size limit. + */ + public void testEvictionTriggersOnRemovalListener() throws Exception { + // Probe for single entry size. + long singleEntryBytes; + try (IndicesBitsetFilterCache probeCache = new IndicesBitsetFilterCache(Settings.EMPTY, threadPool)) { + IndexWriter writer = new IndexWriter( + new ByteBuffersDirectory(), + new IndexWriterConfig(new StandardAnalyzer()).setMergePolicy(new LogByteSizeMergePolicy()) + ); + Document doc = new Document(); + doc.add(new StringField("field", "value", Field.Store.NO)); + writer.addDocument(doc); + writer.commit(); + + DirectoryReader reader = DirectoryReader.open(writer); + reader = OpenSearchDirectoryReader.wrap(reader, new ShardId("probe", "_na_", 0)); + + probeCache.getBitSetProducer(new TermQuery(new Term("field", "value")), NO_OP_LISTENER).getBitSet(reader.leaves().get(0)); + singleEntryBytes = probeCache.getCache().weight(); + + reader.close(); + writer.close(); + } + + // Cache fits only 1 entry. + long cacheSizeBytes = singleEntryBytes; + Settings settings = Settings.builder().put("indices.cache.bitset.size", cacheSizeBytes + "b").build(); + + final AtomicLong removedBytes = new AtomicLong(); + BitsetFilterCache.Listener trackingListener = new BitsetFilterCache.Listener() { + @Override + public void onCache(ShardId shardId, Accountable accountable) {} + + @Override + public void onRemoval(ShardId shardId, Accountable accountable) { + if (accountable != null) { + removedBytes.addAndGet(accountable.ramBytesUsed()); + } + } + }; + + try (IndicesBitsetFilterCache cache = new IndicesBitsetFilterCache(settings, threadPool)) { + IndexWriter writer = new IndexWriter( + new ByteBuffersDirectory(), + new IndexWriterConfig(new StandardAnalyzer()).setMergePolicy(new LogByteSizeMergePolicy()) + ); + + for (int i = 0; i < 2; i++) { + Document doc = new Document(); + doc.add(new StringField("field", "val" + i, Field.Store.NO)); + writer.addDocument(doc); + writer.commit(); + } + + DirectoryReader reader = DirectoryReader.open(writer); + reader = OpenSearchDirectoryReader.wrap(reader, new ShardId("test", "_na_", 0)); + assertThat(reader.leaves().size(), equalTo(2)); + + // Cache first entry. + cache.getBitSetProducer(new TermQuery(new Term("field", "val0")), trackingListener).getBitSet(reader.leaves().get(0)); + assertThat(cache.getCache().count(), equalTo(1)); + assertThat(removedBytes.get(), equalTo(0L)); + + // Cache second entry — should evict the first since limit is 1 entry. + cache.getBitSetProducer(new TermQuery(new Term("field", "val1")), trackingListener).getBitSet(reader.leaves().get(1)); + assertThat(cache.getCache().count(), equalTo(1)); + assertThat(removedBytes.get(), greaterThan(0L)); + + reader.close(); + writer.close(); + } + } + + /** + * Verifies that stale entries from closed readers are purged. + */ + public void testStaleEntriesPurgedAfterReaderClose() throws Exception { + try (IndicesBitsetFilterCache cache = new IndicesBitsetFilterCache(Settings.EMPTY, threadPool)) { + IndexWriter writer = new IndexWriter( + new ByteBuffersDirectory(), + new IndexWriterConfig(new StandardAnalyzer()).setMergePolicy(new LogByteSizeMergePolicy()) + ); + + Document doc = new Document(); + doc.add(new StringField("field", "value", Field.Store.NO)); + writer.addDocument(doc); + writer.commit(); + + DirectoryReader reader = DirectoryReader.open(writer); + reader = OpenSearchDirectoryReader.wrap(reader, new ShardId("test", "_na_", 0)); + + cache.getBitSetProducer(new TermQuery(new Term("field", "value")), NO_OP_LISTENER).getBitSet(reader.leaves().get(0)); + assertThat(cache.getCache().count(), equalTo(1)); + + // Close writer first so no references remain, then close reader. + writer.close(); + reader.close(); + + // Purge stale entries. + cache.purgeStaleEntries(); + assertThat(cache.getCache().count(), equalTo(0)); + } + } + + /** + * Verifies that entries from multiple indices share the same cache and the size limit applies globally. + */ + public void testMultipleIndicesShareCacheWithGlobalSizeLimit() throws Exception { + // Probe for single entry size. + long singleEntryBytes; + try (IndicesBitsetFilterCache probeCache = new IndicesBitsetFilterCache(Settings.EMPTY, threadPool)) { + IndexWriter writer = new IndexWriter( + new ByteBuffersDirectory(), + new IndexWriterConfig(new StandardAnalyzer()).setMergePolicy(new LogByteSizeMergePolicy()) + ); + Document doc = new Document(); + doc.add(new StringField("field", "value", Field.Store.NO)); + writer.addDocument(doc); + writer.commit(); + + DirectoryReader reader = DirectoryReader.open(writer); + reader = OpenSearchDirectoryReader.wrap(reader, new ShardId("probe", "_na_", 0)); + + probeCache.getBitSetProducer(new TermQuery(new Term("field", "value")), NO_OP_LISTENER).getBitSet(reader.leaves().get(0)); + singleEntryBytes = probeCache.getCache().weight(); + + reader.close(); + writer.close(); + } + + // Cache fits 2 entries total across all indices. + long cacheSizeBytes = singleEntryBytes * 2; + Settings settings = Settings.builder().put("indices.cache.bitset.size", cacheSizeBytes + "b").build(); + + try (IndicesBitsetFilterCache cache = new IndicesBitsetFilterCache(settings, threadPool)) { + // Create two separate "indices" (different ShardIds, different writers). + IndexWriter writer1 = new IndexWriter( + new ByteBuffersDirectory(), + new IndexWriterConfig(new StandardAnalyzer()).setMergePolicy(new LogByteSizeMergePolicy()) + ); + Document doc1 = new Document(); + doc1.add(new StringField("field", "val1", Field.Store.NO)); + writer1.addDocument(doc1); + writer1.commit(); + + IndexWriter writer2 = new IndexWriter( + new ByteBuffersDirectory(), + new IndexWriterConfig(new StandardAnalyzer()).setMergePolicy(new LogByteSizeMergePolicy()) + ); + Document doc2 = new Document(); + doc2.add(new StringField("field", "val2", Field.Store.NO)); + writer2.addDocument(doc2); + writer2.commit(); + + IndexWriter writer3 = new IndexWriter( + new ByteBuffersDirectory(), + new IndexWriterConfig(new StandardAnalyzer()).setMergePolicy(new LogByteSizeMergePolicy()) + ); + Document doc3 = new Document(); + doc3.add(new StringField("field", "val3", Field.Store.NO)); + writer3.addDocument(doc3); + writer3.commit(); + + DirectoryReader reader1 = OpenSearchDirectoryReader.wrap(DirectoryReader.open(writer1), new ShardId("index1", "_na_", 0)); + DirectoryReader reader2 = OpenSearchDirectoryReader.wrap(DirectoryReader.open(writer2), new ShardId("index2", "_na_", 0)); + DirectoryReader reader3 = OpenSearchDirectoryReader.wrap(DirectoryReader.open(writer3), new ShardId("index3", "_na_", 0)); + + // Cache one entry from each "index". + cache.getBitSetProducer(new TermQuery(new Term("field", "val1")), NO_OP_LISTENER).getBitSet(reader1.leaves().get(0)); + cache.getBitSetProducer(new TermQuery(new Term("field", "val2")), NO_OP_LISTENER).getBitSet(reader2.leaves().get(0)); + cache.getBitSetProducer(new TermQuery(new Term("field", "val3")), NO_OP_LISTENER).getBitSet(reader3.leaves().get(0)); + + // 3 entries inserted but only 2 fit — global eviction should have kicked in. + assertThat(cache.getCache().count(), lessThanOrEqualTo(2)); + assertThat(cache.getCache().weight(), lessThanOrEqualTo(cacheSizeBytes)); + + reader1.close(); + reader2.close(); + reader3.close(); + writer1.close(); + writer2.close(); + writer3.close(); + } + } +} diff --git a/server/src/test/java/org/opensearch/search/aggregations/bucket/nested/NestedAggregatorTests.java b/server/src/test/java/org/opensearch/search/aggregations/bucket/nested/NestedAggregatorTests.java index 7eb0eedbf89e2..63fc3bbe58b26 100644 --- a/server/src/test/java/org/opensearch/search/aggregations/bucket/nested/NestedAggregatorTests.java +++ b/server/src/test/java/org/opensearch/search/aggregations/bucket/nested/NestedAggregatorTests.java @@ -78,7 +78,6 @@ import org.opensearch.index.query.TermsQueryBuilder; import org.opensearch.index.query.support.NestedScope; import org.opensearch.indices.IndicesBitsetFilterCache; -import org.opensearch.threadpool.TestThreadPool; import org.opensearch.script.MockScriptEngine; import org.opensearch.script.Script; import org.opensearch.script.ScriptEngine; @@ -107,6 +106,7 @@ import org.opensearch.search.aggregations.pipeline.InternalSimpleValue; import org.opensearch.search.aggregations.support.AggregationInspectionHelper; import org.opensearch.search.aggregations.support.ValueType; +import org.opensearch.threadpool.TestThreadPool; import java.io.IOException; import java.util.ArrayList; @@ -1140,8 +1140,9 @@ protected QueryShardContext createQueryShardContext(String fieldName, IndexSetti aggTestIndicesBitsetFilterCache = new IndicesBitsetFilterCache(Settings.EMPTY, aggTestThreadPool); } BitsetFilterCache bitsetFilterCache = new BitsetFilterCache( + indexSettings, aggTestIndicesBitsetFilterCache, - Mockito.mock(IndicesBitsetFilterCache.Listener.class) + Mockito.mock(BitsetFilterCache.Listener.class) ); BitSetProducer nonNestedFilter = bitsetFilterCache.getBitSetProducer(Queries.newNonNestedFilter()); when(queryShardContext.bitsetFilter(any())).thenReturn(nonNestedFilter); diff --git a/server/src/test/java/org/opensearch/search/internal/ContextIndexSearcherTests.java b/server/src/test/java/org/opensearch/search/internal/ContextIndexSearcherTests.java index 6d8dbb5687df3..6ea54e619c277 100644 --- a/server/src/test/java/org/opensearch/search/internal/ContextIndexSearcherTests.java +++ b/server/src/test/java/org/opensearch/search/internal/ContextIndexSearcherTests.java @@ -247,7 +247,7 @@ public void doTestContextIndexSearcher(boolean sparse, boolean deletions) throws w.deleteDocuments(new Term("delete", "yes")); IndexSettings settings = IndexSettingsModule.newIndexSettings("_index", Settings.EMPTY); - IndicesBitsetFilterCache.Listener listener = new IndicesBitsetFilterCache.Listener() { + BitsetFilterCache.Listener listener = new BitsetFilterCache.Listener() { @Override public void onCache(ShardId shardId, Accountable accountable) { @@ -261,7 +261,7 @@ public void onRemoval(ShardId shardId, Accountable accountable) { DirectoryReader reader = OpenSearchDirectoryReader.wrap(DirectoryReader.open(w), new ShardId(settings.getIndex(), 0)); ThreadPool tp = new TestThreadPool("test"); IndicesBitsetFilterCache indicesCache = new IndicesBitsetFilterCache(Settings.EMPTY, tp); - BitsetFilterCache cache = new BitsetFilterCache(indicesCache, listener); + BitsetFilterCache cache = new BitsetFilterCache(settings, indicesCache, listener); Query roleQuery = new TermQuery(new Term("allowed", "yes")); BitSet bitSet = cache.getBitSetProducer(roleQuery).getBitSet(reader.leaves().get(0)); if (sparse) { diff --git a/server/src/test/java/org/opensearch/search/sort/AbstractSortTestCase.java b/server/src/test/java/org/opensearch/search/sort/AbstractSortTestCase.java index 1c0ceb75e3269..c7df3b038fe41 100644 --- a/server/src/test/java/org/opensearch/search/sort/AbstractSortTestCase.java +++ b/server/src/test/java/org/opensearch/search/sort/AbstractSortTestCase.java @@ -210,8 +210,9 @@ protected final QueryShardContext createMockShardContext(IndexSearcher searcher) Settings.builder().put(IndexMetadata.SETTING_VERSION_CREATED, Version.CURRENT).build() ); BitsetFilterCache bitsetFilterCache = new BitsetFilterCache( + idxSettings, Mockito.mock(IndicesBitsetFilterCache.class, Mockito.RETURNS_DEEP_STUBS), - Mockito.mock(IndicesBitsetFilterCache.Listener.class) + Mockito.mock(BitsetFilterCache.Listener.class) ); TriFunction, IndexFieldData> indexFieldDataLookup = ( fieldType, diff --git a/test/framework/src/main/java/org/opensearch/search/aggregations/AggregatorTestCase.java b/test/framework/src/main/java/org/opensearch/search/aggregations/AggregatorTestCase.java index 19d7dfa876f27..14affb7366ac0 100644 --- a/test/framework/src/main/java/org/opensearch/search/aggregations/AggregatorTestCase.java +++ b/test/framework/src/main/java/org/opensearch/search/aggregations/AggregatorTestCase.java @@ -128,8 +128,6 @@ import org.opensearch.index.shard.IndexShard; import org.opensearch.index.shard.SearchOperationListener; import org.opensearch.indices.IndicesBitsetFilterCache; -import org.opensearch.threadpool.TestThreadPool; -import org.opensearch.threadpool.ThreadPool; import org.opensearch.indices.IndicesModule; import org.opensearch.indices.mapper.MapperRegistry; import org.opensearch.plugins.SearchPlugin; @@ -154,6 +152,8 @@ import org.opensearch.search.streaming.FlushMode; import org.opensearch.test.InternalAggregationTestCase; import org.opensearch.test.OpenSearchTestCase; +import org.opensearch.threadpool.TestThreadPool; +import org.opensearch.threadpool.ThreadPool; import org.junit.After; import org.junit.Before; @@ -511,7 +511,7 @@ public boolean shouldCache(Query query) { aggTestIndicesBitsetFilterCache = new IndicesBitsetFilterCache(Settings.EMPTY, aggTestThreadPool); } when(searchContext.bitsetFilterCache()).thenReturn( - new BitsetFilterCache(aggTestIndicesBitsetFilterCache, mock(IndicesBitsetFilterCache.Listener.class)) + new BitsetFilterCache(indexSettings, aggTestIndicesBitsetFilterCache, mock(BitsetFilterCache.Listener.class)) ); IndexShard indexShard = mock(IndexShard.class); when(indexShard.shardId()).thenReturn(new ShardId("test", "test", 0)); diff --git a/test/framework/src/main/java/org/opensearch/test/AbstractBuilderTestCase.java b/test/framework/src/main/java/org/opensearch/test/AbstractBuilderTestCase.java index 0e11a932fb5b7..c7dff0343d34c 100644 --- a/test/framework/src/main/java/org/opensearch/test/AbstractBuilderTestCase.java +++ b/test/framework/src/main/java/org/opensearch/test/AbstractBuilderTestCase.java @@ -453,7 +453,7 @@ private static class ServiceHolder implements Closeable { threadPool ); indicesBitsetFilterCache = new IndicesBitsetFilterCache(Settings.EMPTY, threadPool); - bitsetFilterCache = new BitsetFilterCache(indicesBitsetFilterCache, new IndicesBitsetFilterCache.Listener() { + bitsetFilterCache = new BitsetFilterCache(idxSettings, indicesBitsetFilterCache, new BitsetFilterCache.Listener() { @Override public void onCache(ShardId shardId, Accountable accountable) { From 6c5c8796ac4d0dc1a9ff693039f5def834fb628a Mon Sep 17 00:00:00 2001 From: Sagar Upadhyaya Date: Thu, 9 Apr 2026 16:36:38 -0700 Subject: [PATCH 5/9] Add missing javadoc Signed-off-by: Sagar Upadhyaya --- .../org/opensearch/indices/IndicesBitsetFilterCache.java | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java b/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java index 99e72bc8826e4..74f42a6b04db5 100644 --- a/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java +++ b/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java @@ -228,6 +228,11 @@ public int hashCode() { } } + /** + * Cached value holding the bitset, shard identity, and the per-index listener for stats. + * + * @opensearch.internal + */ @ExperimentalApi public static final class Value { final BitSet bitset; From 66619c5231745105215c9ef2546c453e72d67fbe Mon Sep 17 00:00:00 2001 From: Sagar Upadhyaya Date: Sat, 18 Apr 2026 23:04:27 -0700 Subject: [PATCH 6/9] Address comments and more UTs Signed-off-by: Sagar Upadhyaya --- .../PercolatorQuerySearchTests.java | 6 ++- .../org/opensearch/index/IndexService.java | 4 +- .../index/cache/bitset/BitsetFilterCache.java | 7 ++- .../indices/IndicesBitsetFilterCache.java | 5 +- .../IndicesBitsetFilterCacheTests.java | 48 +++++++++++++++++++ 5 files changed, 62 insertions(+), 8 deletions(-) diff --git a/modules/percolator/src/test/java/org/opensearch/percolator/PercolatorQuerySearchTests.java b/modules/percolator/src/test/java/org/opensearch/percolator/PercolatorQuerySearchTests.java index 5f4925a4ae577..74f346d677f14 100644 --- a/modules/percolator/src/test/java/org/opensearch/percolator/PercolatorQuerySearchTests.java +++ b/modules/percolator/src/test/java/org/opensearch/percolator/PercolatorQuerySearchTests.java @@ -41,13 +41,13 @@ import org.opensearch.core.xcontent.MediaTypeRegistry; import org.opensearch.core.xcontent.XContentBuilder; import org.opensearch.index.IndexService; -import org.opensearch.index.cache.bitset.BitsetFilterCache; import org.opensearch.index.engine.Engine; import org.opensearch.index.fielddata.ScriptDocValues; import org.opensearch.index.query.Operator; import org.opensearch.index.query.QueryBuilder; import org.opensearch.index.query.QueryBuilders; import org.opensearch.index.query.QueryShardContext; +import org.opensearch.indices.IndicesBitsetFilterCache; import org.opensearch.plugins.Plugin; import org.opensearch.script.MockScriptPlugin; import org.opensearch.script.Script; @@ -149,7 +149,9 @@ public void testPercolateQueryWithNestedDocuments_doNotLeakBitsetCacheEntries() .indices() .prepareCreate("test") // to avoid normal document from being cached by BitsetFilterCache - .setSettings(Settings.builder().put(BitsetFilterCache.INDEX_LOAD_RANDOM_ACCESS_FILTERS_EAGERLY_SETTING.getKey(), false)) + .setSettings( + Settings.builder().put(IndicesBitsetFilterCache.INDEX_LOAD_RANDOM_ACCESS_FILTERS_EAGERLY_SETTING.getKey(), false) + ) .setMapping(mapping) ); client().prepareIndex("test") diff --git a/server/src/main/java/org/opensearch/index/IndexService.java b/server/src/main/java/org/opensearch/index/IndexService.java index 9f2d4bb157ac3..aa46b8a5ebffa 100644 --- a/server/src/main/java/org/opensearch/index/IndexService.java +++ b/server/src/main/java/org/opensearch/index/IndexService.java @@ -306,9 +306,7 @@ public IndexService( this.indexSortSupplier = () -> null; } indexFieldData.setListener(new FieldDataCacheListener(this)); - this.bitsetFilterCache = indicesBitsetFilterCache != null - ? new BitsetFilterCache(indexSettings, indicesBitsetFilterCache, new BitsetCacheListener(this)) - : null; + this.bitsetFilterCache = new BitsetFilterCache(indexSettings, indicesBitsetFilterCache, new BitsetCacheListener(this)); this.warmer = new IndexWarmer( threadPool, indexFieldData, diff --git a/server/src/main/java/org/opensearch/index/cache/bitset/BitsetFilterCache.java b/server/src/main/java/org/opensearch/index/cache/bitset/BitsetFilterCache.java index e9527ad666c93..c397e2d554617 100644 --- a/server/src/main/java/org/opensearch/index/cache/bitset/BitsetFilterCache.java +++ b/server/src/main/java/org/opensearch/index/cache/bitset/BitsetFilterCache.java @@ -71,6 +71,10 @@ public final class BitsetFilterCache extends AbstractIndexComponent RemovalListener>, Closeable { + /** + * @deprecated Use {@link IndicesBitsetFilterCache#INDEX_LOAD_RANDOM_ACCESS_FILTERS_EAGERLY_SETTING} instead. + */ + @Deprecated public static final Setting INDEX_LOAD_RANDOM_ACCESS_FILTERS_EAGERLY_SETTING = IndicesBitsetFilterCache.INDEX_LOAD_RANDOM_ACCESS_FILTERS_EAGERLY_SETTING; @@ -123,7 +127,8 @@ public void onClose(IndexReader.CacheKey ownerCoreCacheKey) { @Override public void close() { - // Per-index close is a no-op; the node-level cache manages the lifecycle. + // Per-index close is a no-op; entries are cleaned up by the node-level + // periodic stale-key purge after the index's readers close. } public void clear(String reason) { diff --git a/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java b/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java index 74f42a6b04db5..ac7b2858bfe5a 100644 --- a/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java +++ b/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java @@ -186,14 +186,15 @@ public void purgeStaleEntries() { return; } Set staleSnapshot = new HashSet<>(staleCacheKeys); - staleCacheKeys.removeAll(staleSnapshot); - registeredKeys.removeAll(staleSnapshot); for (BitsetCacheKey key : cache.keys()) { if (staleSnapshot.contains(key.readerCacheKey)) { cache.invalidate(key); } } + + staleCacheKeys.removeAll(staleSnapshot); + registeredKeys.removeAll(staleSnapshot); } public Cache getCache() { diff --git a/server/src/test/java/org/opensearch/indices/IndicesBitsetFilterCacheTests.java b/server/src/test/java/org/opensearch/indices/IndicesBitsetFilterCacheTests.java index a1a6891589973..b6e9cf224dde4 100644 --- a/server/src/test/java/org/opensearch/indices/IndicesBitsetFilterCacheTests.java +++ b/server/src/test/java/org/opensearch/indices/IndicesBitsetFilterCacheTests.java @@ -239,6 +239,54 @@ public void testStaleEntriesPurgedAfterReaderClose() throws Exception { } } + /** + * Verifies that when one index's readers are closed (simulating index close), + * only that index's entries are purged while other indices' entries remain. + */ + public void testIndexCloseOnlyPurgesItsOwnEntries() throws Exception { + try (IndicesBitsetFilterCache cache = new IndicesBitsetFilterCache(Settings.EMPTY, threadPool)) { + // Create two separate "indices" with their own writers. + IndexWriter writer1 = new IndexWriter( + new ByteBuffersDirectory(), + new IndexWriterConfig(new StandardAnalyzer()).setMergePolicy(new LogByteSizeMergePolicy()) + ); + Document doc1 = new Document(); + doc1.add(new StringField("field", "val1", Field.Store.NO)); + writer1.addDocument(doc1); + writer1.commit(); + + IndexWriter writer2 = new IndexWriter( + new ByteBuffersDirectory(), + new IndexWriterConfig(new StandardAnalyzer()).setMergePolicy(new LogByteSizeMergePolicy()) + ); + Document doc2 = new Document(); + doc2.add(new StringField("field", "val2", Field.Store.NO)); + writer2.addDocument(doc2); + writer2.commit(); + + DirectoryReader reader1 = OpenSearchDirectoryReader.wrap(DirectoryReader.open(writer1), new ShardId("index1", "_na_", 0)); + DirectoryReader reader2 = OpenSearchDirectoryReader.wrap(DirectoryReader.open(writer2), new ShardId("index2", "_na_", 0)); + + // Cache one entry from each index. + cache.getBitSetProducer(new TermQuery(new Term("field", "val1")), NO_OP_LISTENER).getBitSet(reader1.leaves().get(0)); + cache.getBitSetProducer(new TermQuery(new Term("field", "val2")), NO_OP_LISTENER).getBitSet(reader2.leaves().get(0)); + assertThat(cache.getCache().count(), equalTo(2)); + + // Simulate index1 close: close its reader. + reader1.close(); + writer1.close(); + cache.purgeStaleEntries(); + + // Only index1's entry should be purged; index2's entry remains. + assertThat(cache.getCache().count(), equalTo(1)); + + reader2.close(); + writer2.close(); + cache.purgeStaleEntries(); + assertThat(cache.getCache().count(), equalTo(0)); + } + } + /** * Verifies that entries from multiple indices share the same cache and the size limit applies globally. */ From 73f28fbe44fdaa745af31016db8324a716feb73f Mon Sep 17 00:00:00 2001 From: Sagar Upadhyaya Date: Sun, 19 Apr 2026 00:06:49 -0700 Subject: [PATCH 7/9] Add javadoc Signed-off-by: Sagar Upadhyaya --- .../java/org/opensearch/indices/IndicesBitsetFilterCache.java | 1 + 1 file changed, 1 insertion(+) diff --git a/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java b/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java index ac7b2858bfe5a..a107810f8c3ea 100644 --- a/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java +++ b/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java @@ -64,6 +64,7 @@ /** * Node-level cache for {@link BitSet} based filters. Manages a single flat cache shared across * all indices on the node, with a configurable size limit and async stale entry cleanup. + * Stale entries from closed readers are purged periodically by a background cleaner task. * * @opensearch.api */ From 94a1f9f85d45dc4be3960b6151b485bd6753bebd Mon Sep 17 00:00:00 2001 From: Sagar Upadhyaya Date: Sun, 19 Apr 2026 01:24:21 -0700 Subject: [PATCH 8/9] Add more UTs for code coverage Signed-off-by: Sagar Upadhyaya --- .../cache/bitset/BitSetFilterCacheTests.java | 78 +++++++++++++++++++ 1 file changed, 78 insertions(+) diff --git a/server/src/test/java/org/opensearch/index/cache/bitset/BitSetFilterCacheTests.java b/server/src/test/java/org/opensearch/index/cache/bitset/BitSetFilterCacheTests.java index f65a2d18b5970..9ab8f3d1c4702 100644 --- a/server/src/test/java/org/opensearch/index/cache/bitset/BitSetFilterCacheTests.java +++ b/server/src/test/java/org/opensearch/index/cache/bitset/BitSetFilterCacheTests.java @@ -67,6 +67,8 @@ import java.util.concurrent.atomic.AtomicLong; import static org.hamcrest.Matchers.equalTo; +import static org.hamcrest.Matchers.notNullValue; +import static org.hamcrest.Matchers.nullValue; public class BitSetFilterCacheTests extends OpenSearchTestCase { @@ -224,6 +226,82 @@ public void testSetNullListener() { } } + public void testDeprecatedConstructorAndCreateListener() throws IOException { + // The deprecated 2-arg constructor sets indicesCache to null. + BitsetFilterCache cache = new BitsetFilterCache(INDEX_SETTINGS, new BitsetFilterCache.Listener() { + @Override + public void onCache(ShardId shardId, Accountable accountable) {} + + @Override + public void onRemoval(ShardId shardId, Accountable accountable) {} + }); + + // createListener returns null when indicesCache is null + assertThat(cache.createListener(threadPool), nullValue()); + + // getBitSetProducer throws when indicesCache is null + expectThrows(IllegalStateException.class, () -> cache.getBitSetProducer(new MatchAllDocsQuery())); + + cache.close(); + } + + public void testCreateListenerWithIndicesCache() throws IOException { + IndicesBitsetFilterCache indicesCache = new IndicesBitsetFilterCache(Settings.EMPTY, threadPool); + BitsetFilterCache cache = new BitsetFilterCache(INDEX_SETTINGS, indicesCache, new BitsetFilterCache.Listener() { + @Override + public void onCache(ShardId shardId, Accountable accountable) {} + + @Override + public void onRemoval(ShardId shardId, Accountable accountable) {} + }); + + // createListener returns non-null when indicesCache is present + assertThat(cache.createListener(threadPool), notNullValue()); + + cache.close(); + indicesCache.close(); + } + + public void testBitsetFromQuery() throws IOException { + Directory dir = newDirectory(); + IndexWriter writer = new IndexWriter(dir, newIndexWriterConfig()); + Document doc = new Document(); + doc.add(new StringField("field", "value", Field.Store.NO)); + writer.addDocument(doc); + writer.commit(); + DirectoryReader reader = DirectoryReader.open(writer); + + // Matching query returns non-null bitset + BitSet bitSet = BitsetFilterCache.bitsetFromQuery(new TermQuery(new Term("field", "value")), reader.leaves().get(0)); + assertNotNull(bitSet); + assertEquals(1, bitSet.cardinality()); + + // Non-matching query returns null bitset + BitSet emptyBitSet = BitsetFilterCache.bitsetFromQuery(new TermQuery(new Term("field", "missing")), reader.leaves().get(0)); + assertNull(emptyBitSet); + + IOUtils.close(reader, writer, dir); + } + + public void testNoOpDelegationMethods() throws IOException { + IndicesBitsetFilterCache indicesCache = new IndicesBitsetFilterCache(Settings.EMPTY, threadPool); + BitsetFilterCache cache = new BitsetFilterCache(INDEX_SETTINGS, indicesCache, new BitsetFilterCache.Listener() { + @Override + public void onCache(ShardId shardId, Accountable accountable) {} + + @Override + public void onRemoval(ShardId shardId, Accountable accountable) {} + }); + + // These are all no-ops delegated to the node-level cache; just verify they don't throw. + cache.onClose(null); + cache.clear("test"); + cache.onRemoval(null); + cache.close(); + + indicesCache.close(); + } + public void testRejectOtherIndex() throws IOException { IndicesBitsetFilterCache indicesCache = new IndicesBitsetFilterCache(Settings.EMPTY, threadPool); BitsetFilterCache cache = new BitsetFilterCache(INDEX_SETTINGS, indicesCache, new BitsetFilterCache.Listener() { From 3af02574d684ae14491941f9e6185ae55fc044c2 Mon Sep 17 00:00:00 2001 From: Sagar Upadhyaya Date: Sun, 19 Apr 2026 01:39:24 -0700 Subject: [PATCH 9/9] Retrigger build Signed-off-by: Sagar Upadhyaya --- .../java/org/opensearch/indices/IndicesBitsetFilterCache.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java b/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java index a107810f8c3ea..0d6830ab75060 100644 --- a/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java +++ b/server/src/main/java/org/opensearch/indices/IndicesBitsetFilterCache.java @@ -64,7 +64,7 @@ /** * Node-level cache for {@link BitSet} based filters. Manages a single flat cache shared across * all indices on the node, with a configurable size limit and async stale entry cleanup. - * Stale entries from closed readers are purged periodically by a background cleaner task. + * Stale entries from closed readers are purged periodically by a background cleanup task. * * @opensearch.api */