From 4ea01d63f788f72ee72df3f94a31d9e43b7ce8bc Mon Sep 17 00:00:00 2001 From: AdityaWaskar Date: Sat, 1 Aug 2026 23:20:29 +0530 Subject: [PATCH 1/2] Fail fast on hybrid query with no resolvable search pipeline Signed-off-by: AdityaWaskar --- .../HybridQuerySearchRequestFilter.java | 51 ++++++- .../neuralsearch/util/HybridQueryUtil.java | 4 + .../NeuralSparseTwoPhaseProcessorTests.java | 4 +- .../HybridQuerySearchRequestFilterTests.java | 138 ++++++++++++++++++ 4 files changed, 189 insertions(+), 8 deletions(-) diff --git a/src/main/java/org/opensearch/neuralsearch/search/HybridQuerySearchRequestFilter.java b/src/main/java/org/opensearch/neuralsearch/search/HybridQuerySearchRequestFilter.java index 8f3b0de02..34f4ce36e 100644 --- a/src/main/java/org/opensearch/neuralsearch/search/HybridQuerySearchRequestFilter.java +++ b/src/main/java/org/opensearch/neuralsearch/search/HybridQuerySearchRequestFilter.java @@ -14,21 +14,31 @@ import org.opensearch.action.search.SearchType; import org.opensearch.action.support.ActionFilter; import org.opensearch.action.support.ActionFilterChain; +import org.opensearch.cluster.metadata.IndexMetadata; import org.opensearch.core.action.ActionResponse; +import org.opensearch.index.IndexSettings; import org.opensearch.index.query.QueryBuilder; import org.opensearch.neuralsearch.query.HybridQueryBuilder; import org.opensearch.neuralsearch.util.HybridQueryUtil; +import org.opensearch.neuralsearch.util.NeuralSearchClusterUtil; +import org.opensearch.search.pipeline.SearchPipelineService; import org.opensearch.tasks.Task; +import java.util.List; + import lombok.extern.log4j.Log4j2; /** - * An ActionFilter that automatically disables batched reduction for hybrid queries. + * An ActionFilter that validates hybrid query requests and disables batched reduction for them. * * This filter intercepts all search requests and checks if they contain a hybrid query. * If a hybrid query is detected with search_type=dfs_query_then_fetch, the request is rejected. - * If a hybrid query is detected, it unconditionally sets batchedReduceSize to Integer.MAX_VALUE - * to disable batched reduction, regardless of any user-specified value. + * If a hybrid query is detected but no search pipeline can be resolved for it (neither inline, + * via the search_pipeline request parameter, nor as every target index's default search + * pipeline), the request is rejected, since without a pipeline the normalization processor + * never runs and the response would otherwise be returned in an unnormalized, internal format. + * Otherwise, it unconditionally sets batchedReduceSize to Integer.MAX_VALUE to disable batched + * reduction, regardless of any user-specified value. * * This prevents the "topDocs already consumed" error that occurs when: * 1. Hybrid query is executed @@ -40,8 +50,6 @@ * The NormalizationProcessor requires access to all shard results simultaneously * to perform score normalization and combination. * - * This filter works transparently without any pipeline or query configuration. - * */ @Log4j2 public class HybridQuerySearchRequestFilter implements ActionFilter { @@ -82,6 +90,10 @@ public indexMetadataList = NeuralSearchClusterUtil.instance().getIndexMetadataList(searchRequest); + return indexMetadataList.isEmpty() == false && indexMetadataList.stream().allMatch(this::hasDefaultSearchPipeline); + } + + private boolean hasDefaultSearchPipeline(IndexMetadata indexMetadata) { + String defaultPipeline = IndexSettings.DEFAULT_SEARCH_PIPELINE.get(indexMetadata.getSettings()); + return isConfiguredPipeline(defaultPipeline); + } + + private boolean isConfiguredPipeline(String pipelineId) { + return Objects.nonNull(pipelineId) + && pipelineId.isBlank() == false + && SearchPipelineService.NOOP_PIPELINE_ID.equals(pipelineId) == false; + } } diff --git a/src/main/java/org/opensearch/neuralsearch/util/HybridQueryUtil.java b/src/main/java/org/opensearch/neuralsearch/util/HybridQueryUtil.java index fad337f2d..55182a423 100644 --- a/src/main/java/org/opensearch/neuralsearch/util/HybridQueryUtil.java +++ b/src/main/java/org/opensearch/neuralsearch/util/HybridQueryUtil.java @@ -34,6 +34,10 @@ public class HybridQueryUtil { public static final String HYBRID_QUERY_DFS_SEARCH_TYPE_NOT_SUPPORTED_MESSAGE = "hybrid query does not support search_type [dfs_query_then_fetch]"; + public static final String HYBRID_QUERY_REQUIRES_SEARCH_PIPELINE_MESSAGE = + "hybrid query requires a search pipeline with a normalization processor to be configured, " + + "either via the search_pipeline request parameter or as the target index's default search pipeline"; + /** * This method validates whether the query object is an instance of hybrid query */ diff --git a/src/test/java/org/opensearch/neuralsearch/processor/NeuralSparseTwoPhaseProcessorTests.java b/src/test/java/org/opensearch/neuralsearch/processor/NeuralSparseTwoPhaseProcessorTests.java index 8c6dbe154..913fee70a 100644 --- a/src/test/java/org/opensearch/neuralsearch/processor/NeuralSparseTwoPhaseProcessorTests.java +++ b/src/test/java/org/opensearch/neuralsearch/processor/NeuralSparseTwoPhaseProcessorTests.java @@ -309,9 +309,7 @@ public void testProcessRequest_whenSortByScoreDescWithTrackScores_thenAddRescore NeuralSparseTwoPhaseProcessor.Factory factory = new NeuralSparseTwoPhaseProcessor.Factory(); NeuralSparseQueryBuilder neuralQueryBuilder = new NeuralSparseQueryBuilder(); SearchRequest searchRequest = new SearchRequest(); - searchRequest.source( - new SearchSourceBuilder().query(neuralQueryBuilder).sort(new ScoreSortBuilder()).trackScores(true) - ); + searchRequest.source(new SearchSourceBuilder().query(neuralQueryBuilder).sort(new ScoreSortBuilder()).trackScores(true)); NeuralSparseTwoPhaseProcessor processor = createTestProcessor(factory, 0.5f, true, 4.0f, 10000); processor.processRequest(searchRequest); assertNotNull(searchRequest.source().rescores()); diff --git a/src/test/java/org/opensearch/neuralsearch/search/HybridQuerySearchRequestFilterTests.java b/src/test/java/org/opensearch/neuralsearch/search/HybridQuerySearchRequestFilterTests.java index 95c4005fa..a1f985331 100644 --- a/src/test/java/org/opensearch/neuralsearch/search/HybridQuerySearchRequestFilterTests.java +++ b/src/test/java/org/opensearch/neuralsearch/search/HybridQuerySearchRequestFilterTests.java @@ -9,31 +9,69 @@ import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; import org.mockito.ArgumentCaptor; +import org.opensearch.Version; import org.opensearch.action.bulk.BulkAction; import org.opensearch.action.bulk.BulkRequest; import org.opensearch.action.search.SearchAction; import org.opensearch.action.search.SearchRequest; import org.opensearch.action.search.SearchType; import org.opensearch.action.support.ActionFilterChain; +import org.opensearch.cluster.ClusterName; +import org.opensearch.cluster.ClusterState; +import org.opensearch.cluster.metadata.IndexMetadata; +import org.opensearch.cluster.metadata.IndexNameExpressionResolver; +import org.opensearch.cluster.metadata.Metadata; +import org.opensearch.cluster.service.ClusterService; +import org.opensearch.common.settings.Settings; +import org.opensearch.common.util.concurrent.ThreadContext; import org.opensearch.core.action.ActionListener; import org.opensearch.core.action.ActionResponse; +import org.opensearch.index.IndexSettings; import org.opensearch.index.query.MatchAllQueryBuilder; import org.opensearch.index.query.MatchQueryBuilder; import org.opensearch.neuralsearch.query.HybridQueryBuilder; import org.opensearch.neuralsearch.query.OpenSearchQueryTestCase; +import org.opensearch.neuralsearch.util.HybridQueryUtil; +import org.opensearch.neuralsearch.util.NeuralSearchClusterUtil; import org.opensearch.search.builder.SearchSourceBuilder; +import org.opensearch.search.pipeline.SearchPipelineService; import org.opensearch.tasks.Task; public class HybridQuerySearchRequestFilterTests extends OpenSearchQueryTestCase { + private static final String TEST_INDEX = "test_index"; + private HybridQuerySearchRequestFilter filter; @Override public void setUp() throws Exception { super.setUp(); filter = new HybridQuerySearchRequestFilter(); + // by default, resolve a default search pipeline for TEST_INDEX so existing tests that + // don't care about pipeline resolution keep exercising the batched-reduce-size behavior + setUpDefaultSearchPipeline(TEST_INDEX, "test-pipeline"); + } + + private void setUpDefaultSearchPipeline(String indexName, String defaultSearchPipelineId) { + Settings.Builder settingsBuilder = Settings.builder() + .put(IndexMetadata.SETTING_VERSION_CREATED, Version.CURRENT) + .put(IndexMetadata.SETTING_NUMBER_OF_SHARDS, 1) + .put(IndexMetadata.SETTING_NUMBER_OF_REPLICAS, 0); + if (defaultSearchPipelineId != null) { + settingsBuilder.put(IndexSettings.DEFAULT_SEARCH_PIPELINE.getKey(), defaultSearchPipelineId); + } + IndexMetadata indexMetadata = IndexMetadata.builder(indexName).settings(settingsBuilder).build(); + Metadata metadata = Metadata.builder().put(indexMetadata, false).build(); + ClusterState clusterState = ClusterState.builder(ClusterName.DEFAULT).metadata(metadata).build(); + + ClusterService clusterService = mock(ClusterService.class); + when(clusterService.state()).thenReturn(clusterState); + + IndexNameExpressionResolver indexNameExpressionResolver = new IndexNameExpressionResolver(new ThreadContext(Settings.EMPTY)); + NeuralSearchClusterUtil.instance().initialize(clusterService, indexNameExpressionResolver); } public void testOrder_thenReturnsZero() { @@ -301,4 +339,104 @@ public void testApply_whenEmptyHybridQuery_thenDisablesBatchedReduction() { assertEquals(Integer.MAX_VALUE, searchRequest.getBatchedReduceSize()); verify(chain).proceed(eq(task), eq(SearchAction.NAME), eq(searchRequest), eq(listener)); } + + @SuppressWarnings("unchecked") + public void testApply_whenHybridQueryWithNoResolvableSearchPipeline_thenFails() { + // no default search pipeline configured on the index, no request-level pipeline set + setUpDefaultSearchPipeline(TEST_INDEX, null); + + HybridQueryBuilder hybridQuery = new HybridQueryBuilder(); + hybridQuery.add(new MatchQueryBuilder("field", "value")); + + SearchRequest searchRequest = new SearchRequest(TEST_INDEX); + SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); + sourceBuilder.query(hybridQuery); + searchRequest.source(sourceBuilder); + + Task task = mock(Task.class); + ActionListener listener = mock(ActionListener.class); + ActionFilterChain chain = mock(ActionFilterChain.class); + + filter.apply(task, SearchAction.NAME, searchRequest, listener, chain); + + ArgumentCaptor exceptionCaptor = ArgumentCaptor.forClass(Exception.class); + verify(listener).onFailure(exceptionCaptor.capture()); + verify(chain, never()).proceed(eq(task), eq(SearchAction.NAME), eq(searchRequest), eq(listener)); + assertTrue(exceptionCaptor.getValue() instanceof IllegalArgumentException); + assertThat(exceptionCaptor.getValue().getMessage(), containsString(HybridQueryUtil.HYBRID_QUERY_REQUIRES_SEARCH_PIPELINE_MESSAGE)); + } + + @SuppressWarnings("unchecked") + public void testApply_whenHybridQueryWithNoopRequestPipelineAndNoIndexDefault_thenFails() { + // request explicitly disables the pipeline (search_pipeline=_none), and index has no default either + setUpDefaultSearchPipeline(TEST_INDEX, null); + + HybridQueryBuilder hybridQuery = new HybridQueryBuilder(); + hybridQuery.add(new MatchQueryBuilder("field", "value")); + + SearchRequest searchRequest = new SearchRequest(TEST_INDEX); + SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); + sourceBuilder.query(hybridQuery); + searchRequest.source(sourceBuilder); + searchRequest.pipeline(SearchPipelineService.NOOP_PIPELINE_ID); + + Task task = mock(Task.class); + ActionListener listener = mock(ActionListener.class); + ActionFilterChain chain = mock(ActionFilterChain.class); + + filter.apply(task, SearchAction.NAME, searchRequest, listener, chain); + + ArgumentCaptor exceptionCaptor = ArgumentCaptor.forClass(Exception.class); + verify(listener).onFailure(exceptionCaptor.capture()); + verify(chain, never()).proceed(eq(task), eq(SearchAction.NAME), eq(searchRequest), eq(listener)); + assertThat(exceptionCaptor.getValue().getMessage(), containsString(HybridQueryUtil.HYBRID_QUERY_REQUIRES_SEARCH_PIPELINE_MESSAGE)); + } + + @SuppressWarnings("unchecked") + public void testApply_whenHybridQueryWithRequestLevelPipelineAndNoIndexDefault_thenProceeds() { + // no default search pipeline on the index, but the request explicitly names one + setUpDefaultSearchPipeline(TEST_INDEX, null); + + HybridQueryBuilder hybridQuery = new HybridQueryBuilder(); + hybridQuery.add(new MatchQueryBuilder("field", "value")); + + SearchRequest searchRequest = new SearchRequest(TEST_INDEX); + SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); + sourceBuilder.query(hybridQuery); + searchRequest.source(sourceBuilder); + searchRequest.pipeline("my-normalization-pipeline"); + + Task task = mock(Task.class); + ActionListener listener = mock(ActionListener.class); + ActionFilterChain chain = mock(ActionFilterChain.class); + + filter.apply(task, SearchAction.NAME, searchRequest, listener, chain); + + verify(listener, never()).onFailure(org.mockito.ArgumentMatchers.any()); + verify(chain).proceed(eq(task), eq(SearchAction.NAME), eq(searchRequest), eq(listener)); + assertEquals(Integer.MAX_VALUE, searchRequest.getBatchedReduceSize()); + } + + @SuppressWarnings("unchecked") + public void testApply_whenHybridQueryWithIndexDefaultSearchPipeline_thenProceeds() { + // index has a default search pipeline configured, request doesn't specify one + setUpDefaultSearchPipeline(TEST_INDEX, "index-default-pipeline"); + + HybridQueryBuilder hybridQuery = new HybridQueryBuilder(); + hybridQuery.add(new MatchQueryBuilder("field", "value")); + + SearchRequest searchRequest = new SearchRequest(TEST_INDEX); + SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); + sourceBuilder.query(hybridQuery); + searchRequest.source(sourceBuilder); + + Task task = mock(Task.class); + ActionListener listener = mock(ActionListener.class); + ActionFilterChain chain = mock(ActionFilterChain.class); + + filter.apply(task, SearchAction.NAME, searchRequest, listener, chain); + + verify(listener, never()).onFailure(org.mockito.ArgumentMatchers.any()); + verify(chain).proceed(eq(task), eq(SearchAction.NAME), eq(searchRequest), eq(listener)); + } } From 023078f0639671f4057d6fdd870a05e47008d38b Mon Sep 17 00:00:00 2001 From: AdityaWaskar Date: Sun, 2 Aug 2026 13:01:10 +0530 Subject: [PATCH 2/2] Add changelog entry for #1922 Signed-off-by: AdityaWaskar --- CHANGELOG.md | 1 + 1 file changed, 1 insertion(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index ee539d9d4..5a78bbe82 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), ### Bug Fixes * [SemanticHighlighter] Fix SemanticHighlighterExtBuilder.toXContent ([#1906](https://github.com/opensearch-project/neural-search/issues/1906)) (query-insights [#651](https://github.com/opensearch-project/query-insights/issues/651)) +* [HybridQuery] Block hybrid query when no search pipeline is configured ([#1922](https://github.com/opensearch-project/neural-search/issues/1922)) ### Infrastructure