From 1ca63ee4d536f86d7febe8fe776d313a8a8722f1 Mon Sep 17 00:00:00 2001 From: Shinsuke Sugaya Date: Tue, 15 Sep 2026 23:17:26 +0900 Subject: [PATCH] fix(rank-fusion): stop waiting without limit for a searcher that does not answer RankFusionProcessor#searchWithMultipleSearchers collected every searcher's result with Future#get() and no timeout. With semantic search enabled, the semantic chunk searcher embeds the query before it searches, so an embedding provider that accepts the connection and never answers held every search for as long as the provider's client kept waiting and retrying, and the search then returned the keyword results with nothing to say that part of it was missing. The main searcher was collected through the same pool, so under load it could also queue for a thread behind searchers that were stuck that way. The searchers other than the main one are now waited for up to rank.fusion.timeout milliseconds (default 10000, the same as query.timeout), counted from their submission. One that has not answered by then is cancelled and left out of that search with a WARN that names it and the setting, and the response is flagged partial and timed out, so the search screen shows labels.process_time_is_exceeded and the search APIs return partial and timed_out. The main searcher now runs on the calling thread: the timeout does not apply to it, since the search engine's own timeouts do, and it no longer needs a pool thread. 0 or less keeps the previous unbounded wait. Only the path where Fess fuses the results itself changes. With rank.fusion.engine.enabled=true the sub-queries are still built on the calling thread before the single fused request is sent. --- .../fess/mylasta/direction/FessConfig.java | 29 +++ .../fess/rank/fusion/RankFusionProcessor.java | 52 ++++- src/main/resources/fess_config.properties | 5 + .../RankFusionProcessorTimeoutTest.java | 212 ++++++++++++++++++ 4 files changed, 289 insertions(+), 9 deletions(-) create mode 100644 src/test/java/org/codelibs/fess/rank/fusion/RankFusionProcessorTimeoutTest.java diff --git a/src/main/java/org/codelibs/fess/mylasta/direction/FessConfig.java b/src/main/java/org/codelibs/fess/mylasta/direction/FessConfig.java index a7308ae06..52a8f7806 100644 --- a/src/main/java/org/codelibs/fess/mylasta/direction/FessConfig.java +++ b/src/main/java/org/codelibs/fess/mylasta/direction/FessConfig.java @@ -1267,6 +1267,9 @@ public interface FessConfig extends FessEnv, org.codelibs.fess.mylasta.direction /** The key of the configuration. e.g. -1 */ String RANK_FUSION_THREADS = "rank.fusion.threads"; + /** The key of the configuration. e.g. 10000 */ + String RANK_FUSION_TIMEOUT = "rank.fusion.timeout"; + /** The key of the configuration. e.g. rf_score */ String RANK_FUSION_score_field = "rank.fusion.score_field"; @@ -6657,6 +6660,23 @@ public interface FessConfig extends FessEnv, org.codelibs.fess.mylasta.direction */ Integer getRankFusionThreadsAsInteger(); + /** + * Get the value for the key 'rank.fusion.timeout'.
+ * The value is, e.g. 10000
+ * comment: Maximum time (milliseconds) to wait for the searchers other than the main one when Fess fuses their results itself. + * @return The value of found property. (NotNull: if not found, exception but basically no way) + */ + String getRankFusionTimeout(); + + /** + * Get the value for the key 'rank.fusion.timeout' as {@link Integer}.
+ * The value is, e.g. 10000
+ * comment: Maximum time (milliseconds) to wait for the searchers other than the main one when Fess fuses their results itself. + * @return The value of found property. (NotNull: if not found, exception but basically no way) + * @throws NumberFormatException When the property is not integer. + */ + Integer getRankFusionTimeoutAsInteger(); + /** * Get the value for the key 'rank.fusion.score_field'.
* The value is, e.g. rf_score
@@ -12517,6 +12537,14 @@ public Integer getRankFusionThreadsAsInteger() { return getAsInteger(FessConfig.RANK_FUSION_THREADS); } + public String getRankFusionTimeout() { + return get(FessConfig.RANK_FUSION_TIMEOUT); + } + + public Integer getRankFusionTimeoutAsInteger() { + return getAsInteger(FessConfig.RANK_FUSION_TIMEOUT); + } + public String getRankFusionScoreField() { return get(FessConfig.RANK_FUSION_score_field); } @@ -14702,6 +14730,7 @@ protected java.util.Map prepareGeneratedDefaultMap() { defaultMap.put(FessConfig.RANK_FUSION_window_size, "200"); defaultMap.put(FessConfig.RANK_FUSION_rank_constant, "20"); defaultMap.put(FessConfig.RANK_FUSION_THREADS, "-1"); + defaultMap.put(FessConfig.RANK_FUSION_TIMEOUT, "10000"); defaultMap.put(FessConfig.RANK_FUSION_score_field, "rf_score"); defaultMap.put(FessConfig.SMB_ROLE_FROM_FILE, "true"); defaultMap.put(FessConfig.SMB_AVAILABLE_SID_TYPES, "1,2,4:2,5:1"); diff --git a/src/main/java/org/codelibs/fess/rank/fusion/RankFusionProcessor.java b/src/main/java/org/codelibs/fess/rank/fusion/RankFusionProcessor.java index 56aca41b8..098b3a056 100644 --- a/src/main/java/org/codelibs/fess/rank/fusion/RankFusionProcessor.java +++ b/src/main/java/org/codelibs/fess/rank/fusion/RankFusionProcessor.java @@ -32,6 +32,7 @@ import java.util.concurrent.Executors; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import java.util.stream.Collectors; import org.apache.logging.log4j.LogManager; @@ -325,6 +326,12 @@ protected RankFusionSearcher[] getAvailableSearchers() { * Executes searches in parallel across all provided searchers, then combines the results * using rank fusion algorithms to produce a unified result set. * + *

The main searcher runs on the calling thread and is always waited for. The other searchers + * run on the executor and are waited for up to {@code rank.fusion.timeout} milliseconds, counted + * from their submission. One that has not answered by then - an embedding provider that accepted + * the query and never replied, for instance - is left out of this search, and the response is + * flagged as partial and timed out.

+ * * @param searchers array of searchers to use for concurrent searching * @param query the search query string * @param params search request parameters including pagination and filters @@ -380,9 +387,12 @@ protected List> searchWithMultipleSearchers(final RankFusion if (logger.isDebugEnabled()) { logger.debug("Search parameters: windowSize={}, sizePerSearcher={}, rankConstant={}", windowSize, size, rankConstant); } + final Integer timeout = fessConfig.getRankFusionTimeoutAsInteger(); + final long timeoutNanos = timeout != null && timeout > 0 ? TimeUnit.MILLISECONDS.toNanos(timeout) : 0L; + final long startTime = System.nanoTime(); final List> resultList = new ArrayList<>(); - for (int i = 0; i < searchers.length; i++) { - final SearchRequestParams reqParams = new SearchRequestParamsWrapper(params, 0, i == 0 ? windowSize : size); + for (int i = 1; i < searchers.length; i++) { + final SearchRequestParams reqParams = new SearchRequestParamsWrapper(params, 0, size); final RankFusionSearcher searcher = searchers[i]; resultList.add(executorService.submit(() -> { try { @@ -399,13 +409,37 @@ protected List> searchWithMultipleSearchers(final RankFusion } })); } - final SearchResult[] results = resultList.stream().map(future -> { + final SearchResult[] results = new SearchResult[searchers.length]; + // The main searcher runs on the calling thread. rank.fusion.timeout does not apply to it - + // the search engine's own timeouts do - and it never queues for a pool thread behind + // searchers that are still waiting on an unresponsive provider. + try { + results[0] = searchers[0].search(query, new SearchRequestParamsWrapper(params, 0, windowSize), userBean); + } catch (final InvalidQueryException | ResultOffsetExceededException | InvalidAccessTokenException e) { + throw e; + } catch (final Exception e) { + logger.warn("Search operation failed with exception", e); + results[0] = SearchResult.create().build(); + } + boolean searcherTimedOut = false; + for (int i = 1; i < searchers.length; i++) { + final Future future = resultList.get(i - 1); try { - return future.get(); + if (timeoutNanos > 0L) { + results[i] = future.get(Math.max(0L, timeoutNanos - (System.nanoTime() - startTime)), TimeUnit.NANOSECONDS); + } else { + results[i] = future.get(); + } + } catch (final TimeoutException e) { + future.cancel(true); + logger.warn("{} did not return results within {}ms ({}); it is left out of this search. query={}", searchers[i].getName(), + timeout, FessConfig.RANK_FUSION_TIMEOUT, query); + searcherTimedOut = true; + results[i] = SearchResult.create().build(); } catch (final InterruptedException e) { logger.warn("Search operation was interrupted", e); Thread.currentThread().interrupt(); // Restore interrupt status - return SearchResult.create().build(); + results[i] = SearchResult.create().build(); } catch (final ExecutionException e) { if (e.getCause() instanceof final InvalidQueryException iqe) { throw iqe; @@ -417,9 +451,9 @@ protected List> searchWithMultipleSearchers(final RankFusion throw iate; } logger.warn("Search operation failed with exception", e.getCause()); - return SearchResult.create().build(); + results[i] = SearchResult.create().build(); } - }).toArray(SearchResult[]::new); + } final String scoreField = fessConfig.getRankFusionScoreField(); final Map> documentsByIdMap = new HashMap<>(); @@ -489,8 +523,8 @@ protected List> searchWithMultipleSearchers(final RankFusion allRecordCount += offset; } return createResponseList(extractList(fusedDocs, pageSize, startPosition), allRecordCount, mainResult.getAllRecordCountRelation(), - mainResult.getQueryTime(), mainResult.isPartialResults(), mainResult.isTimedOut(), mainResult.isShardFailed(), - mainResult.getFacetResponse(), startPosition, pageSize, offset); + mainResult.getQueryTime(), mainResult.isPartialResults() || searcherTimedOut, mainResult.isTimedOut() || searcherTimedOut, + mainResult.isShardFailed(), mainResult.getFacetResponse(), startPosition, pageSize, offset); } /** diff --git a/src/main/resources/fess_config.properties b/src/main/resources/fess_config.properties index 84b5f0d23..7f6c597a2 100644 --- a/src/main/resources/fess_config.properties +++ b/src/main/resources/fess_config.properties @@ -1042,6 +1042,11 @@ rank.fusion.window_size=200 rank.fusion.rank_constant=20 # Number of threads for rank fusion. rank.fusion.threads=-1 +# Maximum time (milliseconds) to wait for the searchers other than the main one when Fess fuses +# their results itself (rank.fusion.engine.enabled=false). A searcher that has not answered by +# then is left out of that search, and the results are flagged as partial and timed out. The main +# searcher is always waited for. 0 or less waits without a limit. +rank.fusion.timeout=10000 # Score field for rank fusion. rank.fusion.score_field=rf_score # Whether the search engine performs rank fusion. When true, the searchers that can take diff --git a/src/test/java/org/codelibs/fess/rank/fusion/RankFusionProcessorTimeoutTest.java b/src/test/java/org/codelibs/fess/rank/fusion/RankFusionProcessorTimeoutTest.java new file mode 100644 index 000000000..35c69b6cd --- /dev/null +++ b/src/test/java/org/codelibs/fess/rank/fusion/RankFusionProcessorTimeoutTest.java @@ -0,0 +1,212 @@ +/* + * Copyright 2012-2025 CodeLibs Project and the Others. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, + * either express or implied. See the License for the specific language + * governing permissions and limitations under the License. + */ +package org.codelibs.fess.rank.fusion; + +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +import org.apache.lucene.search.TotalHits.Relation; +import org.codelibs.fess.entity.SearchRequestParams; +import org.codelibs.fess.mylasta.action.FessUserBean; +import org.codelibs.fess.mylasta.direction.FessConfig; +import org.codelibs.fess.rank.fusion.RankFusionProcessorErrorHandlingTest.SlowSearcher; +import org.codelibs.fess.rank.fusion.RankFusionProcessorErrorHandlingTest.TestSearchRequestParams; +import org.codelibs.fess.rank.fusion.RankFusionProcessorErrorHandlingTest.TestSearcher; +import org.codelibs.fess.rank.fusion.SearchResult.SearchResultBuilder; +import org.codelibs.fess.unit.UnitFessTestCase; +import org.codelibs.fess.util.ComponentUtil; +import org.codelibs.fess.util.QueryResponseList; +import org.dbflute.optional.OptionalThing; +import org.junit.jupiter.api.Test; + +/** + * How long rank fusion waits for its searchers ({@code rank.fusion.timeout}). + */ +public class RankFusionProcessorTimeoutTest extends UnitFessTestCase { + + private static final String ID_FIELD = "_id"; + + /** + * A searcher that does not answer - an embedding provider that accepted the query and never + * replied - must not hold the search: the main searcher's results come back once the timeout + * elapses, flagged as partial and timed out. + */ + @Test + public void test_searcherThatDoesNotAnswer_isLeftOutAfterTimeout() throws Exception { + givenTimeout(200); + try (RankFusionProcessor processor = new RankFusionProcessor()) { + processor.setSearcher(new TestSearcher(100)); + final HangingSearcher hanging = new HangingSearcher(5000); + processor.register(hanging); + processor.init(); + + final long start = System.currentTimeMillis(); + final List> results = processor.search("*", new TestSearchRequestParams(0, 10, 0), OptionalThing.empty()); + final long elapsed = System.currentTimeMillis() - start; + + assertTrue("search waited " + elapsed + "ms for a searcher that does not answer", elapsed < 3000); + assertEquals(10, results.size()); + for (final Map doc : results) { + assertFalse("unexpected document from the hanging searcher: " + doc, doc.get(ID_FIELD).toString().startsWith("hang_")); + } + final QueryResponseList response = (QueryResponseList) results; + assertTrue("results without a searcher must be flagged partial", response.isPartialResults()); + assertTrue("a searcher left out for the timeout must be reported as a timeout", response.isTimedOut()); + assertFalse(response.isShardFailed()); + } + } + + /** + * A searcher that answers within the timeout is fused as before. + */ + @Test + public void test_searcherWithinTimeout_isFused() throws Exception { + givenTimeout(5000); + try (RankFusionProcessor processor = new RankFusionProcessor()) { + processor.setSearcher(new TestSearcher(100)); + processor.register(new SlowSearcher(50)); + processor.init(); + + final List> results = processor.search("*", new TestSearchRequestParams(0, 10, 0), OptionalThing.empty()); + + assertTrue("the searcher that answered in time is missing: " + results, containsIdPrefix(results, "slow_")); + final QueryResponseList response = (QueryResponseList) results; + assertFalse(response.isPartialResults()); + assertFalse(response.isTimedOut()); + } + } + + /** + * The timeout bounds the other searchers, not the main one: a slow main searcher is still + * waited for, and the other searchers that finished meanwhile are fused with it. + */ + @Test + public void test_mainSearcherIsNotCutByTimeout() throws Exception { + givenTimeout(100); + try (RankFusionProcessor processor = new RankFusionProcessor()) { + processor.setSearcher(new SlowSearcher(500)); + processor.register(new TestSearcher(100)); + processor.init(); + + final List> results = processor.search("*", new TestSearchRequestParams(0, 10, 0), OptionalThing.empty()); + + assertTrue("the main searcher's results are missing: " + results, containsIdPrefix(results, "slow_")); + final QueryResponseList response = (QueryResponseList) results; + assertFalse(response.isPartialResults()); + assertFalse(response.isTimedOut()); + } + } + + /** + * Zero or less waits for every searcher without a limit. + */ + @Test + public void test_nonPositiveTimeout_waitsForEverySearcher() throws Exception { + givenTimeout(0); + try (RankFusionProcessor processor = new RankFusionProcessor()) { + processor.setSearcher(new TestSearcher(100)); + processor.register(new SlowSearcher(300)); + processor.init(); + + final List> results = processor.search("*", new TestSearchRequestParams(0, 10, 0), OptionalThing.empty()); + + assertTrue("the slow searcher's results are missing: " + results, containsIdPrefix(results, "slow_")); + final QueryResponseList response = (QueryResponseList) results; + assertFalse(response.isPartialResults()); + assertFalse(response.isTimedOut()); + } + } + + private static boolean containsIdPrefix(final List> results, final String prefix) { + return results.stream().anyMatch(doc -> doc.get(ID_FIELD).toString().startsWith(prefix)); + } + + private void givenTimeout(final int timeoutMillis) { + ComponentUtil.setFessConfig(new FessConfig.SimpleImpl() { + private static final long serialVersionUID = 1L; + + @Override + public Integer getRankFusionTimeoutAsInteger() { + return Integer.valueOf(timeoutMillis); + } + + @Override + public Integer getPagingSearchPageMaxSizeAsInteger() { + return Integer.valueOf(100); + } + + @Override + public Integer getRankFusionWindowSizeAsInteger() { + return Integer.valueOf(200); + } + + @Override + public Integer getRankFusionRankConstantAsInteger() { + return Integer.valueOf(20); + } + + @Override + public Integer getRankFusionThreadsAsInteger() { + return Integer.valueOf(-1); + } + + @Override + public String getRankFusionScoreField() { + return "rf_score"; + } + + @Override + public boolean isRankFusionEngineEnabled() { + return false; + } + + @Override + public String getIndexFieldId() { + return ID_FIELD; + } + }); + } + + /** + * Searcher that blocks like a provider that accepted the request and never answered. + */ + static class HangingSearcher extends RankFusionSearcher { + private final long blockMs; + + HangingSearcher(final long blockMs) { + this.blockMs = blockMs; + } + + @Override + protected SearchResult search(final String query, final SearchRequestParams params, final OptionalThing userBean) { + try { + Thread.sleep(blockMs); + } catch (final InterruptedException e) { + Thread.currentThread().interrupt(); + } + final SearchResultBuilder builder = SearchResult.create(); + for (int i = 0; i < 10; i++) { + final Map doc = new HashMap<>(); + doc.put(ID_FIELD, "hang_" + i); + builder.addDocument(doc); + } + builder.allRecordCount(10); + builder.allRecordCountRelation(Relation.EQUAL_TO.toString()); + return builder.build(); + } + } +}