diff --git a/common/src/main/java/org/opensearch/sql/common/setting/Settings.java b/common/src/main/java/org/opensearch/sql/common/setting/Settings.java index 67643a80add..bf686a205fc 100644 --- a/common/src/main/java/org/opensearch/sql/common/setting/Settings.java +++ b/common/src/main/java/org/opensearch/sql/common/setting/Settings.java @@ -48,7 +48,6 @@ public enum Key { /** Query Settings. */ FIELD_TYPE_TOLERANCE("plugins.query.field_type_tolerance"), - PARTIAL_RESULT_ON_MAPPING_CONFLICT("plugins.query.partial_result.on_mapping_conflict.enabled"), /** Common Settings for SQL and PPL. */ QUERY_MEMORY_LIMIT("plugins.query.memory_limit"), diff --git a/core/src/main/java/org/opensearch/sql/calcite/CalcitePlanContext.java b/core/src/main/java/org/opensearch/sql/calcite/CalcitePlanContext.java index 7ba2c1df582..0fbefa4a17a 100644 --- a/core/src/main/java/org/opensearch/sql/calcite/CalcitePlanContext.java +++ b/core/src/main/java/org/opensearch/sql/calcite/CalcitePlanContext.java @@ -65,32 +65,13 @@ public class CalcitePlanContext { public static final ThreadLocal executionPool = new ThreadLocal<>(); /** - * Non-fatal warnings raised during planning (e.g. a partial result over a subset of indices) to - * be attached to the query response by the execution engine. Drained in {@code - * OpenSearchExecutionEngine.buildResultSet} and cleared with the other lifecycle signals so it - * never leaks onto the next query on a pooled worker thread. + * Non-fatal warnings raised while planning or executing a query, to be attached to its response. + * Drained in {@code OpenSearchExecutionEngine.buildResultSet} and cleared with the other + * lifecycle signals so it never leaks onto the next query on a pooled worker thread. */ private static final ThreadLocal> pendingWarnings = ThreadLocal.withInitial(ArrayList::new); - /** - * Whether the current query's response format can carry a warnings channel. Set on the worker - * thread from the plan (see {@code QueryPlan#execute}) rather than the transport thread, so the - * partial-result gate survives the transport→worker handoff — the security plugin's interceptor - * drops Log4j {@code ThreadContext}, which is where this used to live. Cleared per query. - */ - private static final ThreadLocal warningsSupported = - ThreadLocal.withInitial(() -> false); - - /** - * Per-request partial-result override, carried off Log4j {@code ThreadContext} onto the plan (see - * {@code QueryPlan#execute}) for the same reason as {@link #warningsSupported}: the security - * plugin's interceptor drops {@code ThreadContext} on the transport→worker handoff. {@code null} - * defers to the cluster setting; {@code true}/{@code false} force partial mode on/off for this - * query. Cleared per query. - */ - private static final ThreadLocal partialResultOverride = new ThreadLocal<>(); - /** Thread-local switch that tells whether the current query prefers legacy behavior. */ private static final ThreadLocal legacyPreferredFlag = ThreadLocal.withInitial(() -> true); @@ -279,8 +260,6 @@ public static void clearTimewrapSignals() { timewrapSeries.set(null); executionPool.set(null); pendingWarnings.remove(); - warningsSupported.set(false); - partialResultOverride.remove(); } /** Records a non-fatal warning to be attached to the response for the current query. */ @@ -288,35 +267,6 @@ public static void addWarning(Warning warning) { pendingWarnings.get().add(warning); } - /** Records whether the current query's response format can surface warnings. */ - public static void setWarningsSupported(boolean supported) { - warningsSupported.set(supported); - } - - /** - * @return whether the current query's response format can surface warnings; false when unset, so - * a caller that never declared support cannot get a silent partial result. - */ - public static boolean isWarningsSupported() { - return warningsSupported.get(); - } - - /** - * Records the per-request partial-result override for the current query. {@code null} defers to - * the cluster setting; {@code true}/{@code false} force partial mode on/off. - */ - public static void setPartialResultOverride(Boolean override) { - partialResultOverride.set(override); - } - - /** - * @return the per-request partial-result override, or {@code null} to defer to the cluster - * setting. - */ - public static Boolean getPartialResultOverride() { - return partialResultOverride.get(); - } - /** * Returns and clears the warnings collected for the current query, de-duplicated by value. The * planner may fire a rule that raises a warning more than once for equivalent plan alternatives, @@ -342,24 +292,18 @@ public static class ThreadLocalSnapshot { final String timewrapUnitName; final String timewrapSeries; final String executionPool; - final boolean warningsSupported; - final Boolean partialResultOverride; private ThreadLocalSnapshot( boolean skipEncoding, boolean stripNullColumns, String timewrapUnitName, String timewrapSeries, - String executionPool, - boolean warningsSupported, - Boolean partialResultOverride) { + String executionPool) { this.skipEncoding = skipEncoding; this.stripNullColumns = stripNullColumns; this.timewrapUnitName = timewrapUnitName; this.timewrapSeries = timewrapSeries; this.executionPool = executionPool; - this.warningsSupported = warningsSupported; - this.partialResultOverride = partialResultOverride; } } @@ -370,9 +314,7 @@ public static ThreadLocalSnapshot snapshotThreadLocals() { stripNullColumns.get(), timewrapUnitName.get(), timewrapSeries.get(), - executionPool.get(), - warningsSupported.get(), - partialResultOverride.get()); + executionPool.get()); } /** Restore thread-local state from a snapshot. */ @@ -382,8 +324,6 @@ public static void restoreThreadLocals(ThreadLocalSnapshot snapshot) { timewrapUnitName.set(snapshot.timewrapUnitName); timewrapSeries.set(snapshot.timewrapSeries); executionPool.set(snapshot.executionPool); - warningsSupported.set(snapshot.warningsSupported); - partialResultOverride.set(snapshot.partialResultOverride); } public void pushForeachBindings( diff --git a/core/src/main/java/org/opensearch/sql/executor/Warning.java b/core/src/main/java/org/opensearch/sql/executor/Warning.java index daa64c2f118..479645d4355 100644 --- a/core/src/main/java/org/opensearch/sql/executor/Warning.java +++ b/core/src/main/java/org/opensearch/sql/executor/Warning.java @@ -9,17 +9,16 @@ /** * A non-fatal notice attached to an otherwise-successful query response. Carried through the - * response path so consumers can distinguish a correct-but-noteworthy result (e.g. a partial result - * over a subset of indices) from a plain success, without turning it into an error. + * response path so consumers can distinguish a correct-but-noteworthy result from a plain success, + * without turning it into an error. */ @Data public class Warning { /** - * The result is complete for the indices it covers but omits one or more indices that could not - * be served (e.g. a mapping conflict that prevents aggregation pushdown). This is a cross-surface - * contract: consumers such as OpenSearch Dashboards branch on this {@code type} value, so it must - * not change without coordinating those consumers. + * The response does not cover everything the query asked for. This is a cross-surface contract: + * consumers such as OpenSearch Dashboards branch on this {@code type} value, so it must not + * change without coordinating those consumers. */ public static final String TYPE_PARTIAL_RESULT = "PARTIAL_RESULT"; diff --git a/core/src/main/java/org/opensearch/sql/executor/execution/AbstractPlan.java b/core/src/main/java/org/opensearch/sql/executor/execution/AbstractPlan.java index 1e8835b4815..fbdabe2fa44 100644 --- a/core/src/main/java/org/opensearch/sql/executor/execution/AbstractPlan.java +++ b/core/src/main/java/org/opensearch/sql/executor/execution/AbstractPlan.java @@ -7,7 +7,6 @@ import lombok.Getter; import lombok.RequiredArgsConstructor; -import lombok.Setter; import org.opensearch.sql.ast.statement.ExplainMode; import org.opensearch.sql.common.response.ResponseListener; import org.opensearch.sql.executor.ExecutionEngine; @@ -24,22 +23,6 @@ public abstract class AbstractPlan { @Getter protected final QueryType queryType; - /** - * Whether the response format can carry a warnings channel. Set from the request on the transport - * thread and read on the worker (see {@code QueryPlan#execute}), so a feature that returns a - * partial result never silently drops data into a warnings-incapable format across a thread - * handoff. Defaults to false. - */ - @Getter @Setter private boolean warningsSupported = false; - - /** - * Per-request partial-result override, carried from the request the same way as {@link - * #warningsSupported}. {@code null} defers to the cluster setting; {@code true}/{@code false} - * force partial mode on/off. Set on the transport thread, applied on the worker (see {@code - * QueryPlan#execute}) so it survives the security transport→worker handoff. - */ - @Getter @Setter private Boolean partialResultOverride = null; - /** Start query execution. */ public abstract void execute(); diff --git a/core/src/main/java/org/opensearch/sql/executor/execution/QueryPlan.java b/core/src/main/java/org/opensearch/sql/executor/execution/QueryPlan.java index 570bf21f97f..d78036f0faa 100644 --- a/core/src/main/java/org/opensearch/sql/executor/execution/QueryPlan.java +++ b/core/src/main/java/org/opensearch/sql/executor/execution/QueryPlan.java @@ -11,7 +11,6 @@ import org.opensearch.sql.ast.tree.HighlightConfig; import org.opensearch.sql.ast.tree.Paginate; import org.opensearch.sql.ast.tree.UnresolvedPlan; -import org.opensearch.sql.calcite.CalcitePlanContext; import org.opensearch.sql.common.response.ResponseListener; import org.opensearch.sql.executor.ExecutionEngine; import org.opensearch.sql.executor.QueryId; @@ -106,11 +105,6 @@ public QueryPlan( @Override public void execute() { - // Runs on the worker thread; carry warnings support and the per-request partial-result override - // from the request off the plan so the partial-result gate reads them without depending on - // Log4j ThreadContext (dropped under security on the transport→worker handoff). - CalcitePlanContext.setWarningsSupported(isWarningsSupported()); - CalcitePlanContext.setPartialResultOverride(getPartialResultOverride()); if (pageSize.isPresent()) { queryService.execute( new Paginate(pageSize.get(), plan), diff --git a/core/src/test/java/org/opensearch/sql/executor/execution/QueryPlanTest.java b/core/src/test/java/org/opensearch/sql/executor/execution/QueryPlanTest.java index 0080f17c177..3220c5d28f7 100644 --- a/core/src/test/java/org/opensearch/sql/executor/execution/QueryPlanTest.java +++ b/core/src/test/java/org/opensearch/sql/executor/execution/QueryPlanTest.java @@ -5,7 +5,6 @@ package org.opensearch.sql.executor.execution; -import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assertions.fail; @@ -17,7 +16,6 @@ import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; -import java.util.concurrent.atomic.AtomicBoolean; import org.apache.commons.lang3.NotImplementedException; import org.junit.jupiter.api.DisplayNameGeneration; import org.junit.jupiter.api.DisplayNameGenerator; @@ -27,7 +25,6 @@ import org.mockito.junit.jupiter.MockitoExtension; import org.opensearch.sql.ast.statement.ExplainMode; import org.opensearch.sql.ast.tree.UnresolvedPlan; -import org.opensearch.sql.calcite.CalcitePlanContext; import org.opensearch.sql.common.response.ResponseListener; import org.opensearch.sql.executor.DefaultExecutionEngine; import org.opensearch.sql.executor.ExecutionEngine; @@ -61,53 +58,6 @@ public void execute_no_page_size() { verify(queryService, times(1)).execute(any(), any(), any(), anyBoolean(), any()); } - @Test - public void warnings_supported_flag_reaches_execution_thread() throws InterruptedException { - QueryPlan query = new QueryPlan(queryId, queryType, plan, queryService, queryListener); - query.setWarningsSupported(true); - - // Configure the plan here but run execute() on another thread, mirroring the transport->worker - // handoff. The flag must ride the plan object, not a thread-local the handoff can drop. - AtomicBoolean defaultedBeforeExecute = new AtomicBoolean(true); - AtomicBoolean seenOnWorker = new AtomicBoolean(false); - Thread worker = - new Thread( - () -> { - defaultedBeforeExecute.set(CalcitePlanContext.isWarningsSupported()); - query.execute(); - seenOnWorker.set(CalcitePlanContext.isWarningsSupported()); - }); - worker.start(); - worker.join(); - - assertFalse( - defaultedBeforeExecute.get(), "worker thread should default to no warnings support"); - assertTrue(seenOnWorker.get(), "execute() must carry warningsSupported onto the worker thread"); - verify(queryService, times(1)).execute(any(), any(), any(), anyBoolean(), any()); - } - - @Test - public void warnings_unsupported_plan_resets_flag_on_reused_worker_thread() - throws InterruptedException { - // Plan defaults to warningsSupported=false. - QueryPlan query = new QueryPlan(queryId, queryType, plan, queryService, queryListener); - - AtomicBoolean seenOnWorker = new AtomicBoolean(true); - Thread worker = - new Thread( - () -> { - // Simulate a pooled worker left "supported" by a prior query. - CalcitePlanContext.setWarningsSupported(true); - query.execute(); - seenOnWorker.set(CalcitePlanContext.isWarningsSupported()); - }); - worker.start(); - worker.join(); - - assertFalse( - seenOnWorker.get(), "a warnings-unsupported plan must reset the flag on a reused thread"); - } - @Test public void explain_no_page_size() { QueryPlan query = new QueryPlan(queryId, queryType, plan, queryService, queryListener); diff --git a/docs/user/admin/settings.rst b/docs/user/admin/settings.rst index e91b6bdad2f..f6e9b16c6c0 100644 --- a/docs/user/admin/settings.rst +++ b/docs/user/admin/settings.rst @@ -358,53 +358,6 @@ Result set:: } } -plugins.query.partial_result.on_mapping_conflict.enabled [Experimental] -======================================================================= - -Version -------- -Since 3.9 - -Description ------------ - -This setting is experimental; its name, values, and default may change in a future release. Controls how an aggregation behaves when its group-by field is mapped inconsistently across the queried indices -- for example ``keyword`` in some indices of a wildcard pattern and ``text`` (without a ``.keyword`` sub-field) in others. Such a field collapses to ``text``-without-``.keyword`` across the pattern, which has no doc values, so the aggregation cannot be pushed down natively and instead runs as a per-document script over ``_source`` -- correct, but a full scan of every document. - -When this setting is ``false`` (the default), that complete-but-slow result is returned. When set to ``true``, the aggregation is pushed down over only the subset of indices where the field is aggregatable, and the response carries a ``PARTIAL_RESULT`` warning naming the excluded indices and the remedy (map the field as ``keyword`` everywhere). The result is therefore **partial** -- documents in the excluded indices are not counted -- so the setting is off by default and only takes effect for response formats that can surface the warning (the JSON format; CSV/raw/visualization responses fall through to the complete result rather than silently dropping data). - -The behavior can also be overridden per request with the ``partial_result`` boolean field in the query body, which takes precedence over this cluster setting. Here is an example enabling it at the cluster level:: - - >> curl -H 'Content-Type: application/json' -X PUT localhost:9200/_plugins/_query/settings -d '{ - "transient" : { - "plugins.query.partial_result.on_mapping_conflict.enabled" : true - } - }' - -Result set:: - - { - "acknowledged" : true, - "persistent" : { }, - "transient" : { - "plugins" : { - "query" : { - "partial_result" : { - "on_mapping_conflict" : { - "enabled" : "true" - } - } - } - } - } - } - -Per-request override example, opting a single query into a partial result regardless of the cluster setting:: - - >> curl -H 'Content-Type: application/json' -X POST localhost:9200/_plugins/_ppl -d '{ - "query" : "source=logs-* | stats count() by service", - "partial_result" : true - }' - plugins.query.buckets ===================== diff --git a/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalcitePartialResultOnMappingConflictIT.java b/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalcitePartialResultOnMappingConflictIT.java deleted file mode 100644 index 929db438c45..00000000000 --- a/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalcitePartialResultOnMappingConflictIT.java +++ /dev/null @@ -1,463 +0,0 @@ -/* - * Copyright OpenSearch Contributors - * SPDX-License-Identifier: Apache-2.0 - */ - -package org.opensearch.sql.calcite.remote; - -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertTrue; -import static org.opensearch.sql.util.Capability.CROSS_INDEX_INCOMPATIBLE_TYPES; -import static org.opensearch.sql.util.MatcherUtils.rows; -import static org.opensearch.sql.util.MatcherUtils.verifyDataRows; -import static org.opensearch.sql.util.TestUtils.createIndexByRestClient; -import static org.opensearch.sql.util.TestUtils.isIndexExist; -import static org.opensearch.sql.util.TestUtils.performRequest; - -import java.io.IOException; -import org.json.JSONArray; -import org.json.JSONObject; -import org.junit.After; -import org.junit.Test; -import org.opensearch.client.Request; -import org.opensearch.sql.common.setting.Settings; -import org.opensearch.sql.ppl.PPLIntegTestCase; -import org.opensearch.sql.util.RequiresCapability; - -/** - * End-to-end tests for the partial-result path on a text/keyword mapping conflict. A field mapped - * as {@code keyword} in one index and {@code text} (without a {@code .keyword} sub-field) in - * another collapses to text-without-keyword across the wildcard pattern. Aggregating that field is - * possible but expensive: since #5646 it pushes down as a per-document {@code _source} script that - * reads every document. - * - *

When {@code plugins.query.partial_result.on_mapping_conflict.enabled} is on, the aggregation - * is instead pushed down natively over just the aggregatable (keyword) index subset — far faster, - * but incomplete — and the response carries a {@code PARTIAL_RESULT} warning naming the - * excluded indices. When off, the complete (slow) result is returned with no warning. - */ -public class CalcitePartialResultOnMappingConflictIT extends PPLIntegTestCase { - - private static final String KEYWORD_INDEX = "partial_conflict_keyword"; - private static final String TEXT_INDEX = "partial_conflict_text"; - private static final String PATTERN = "partial_conflict_*"; - - private static final String NESTED_KEYWORD_INDEX = "partial_nested_keyword"; - private static final String NESTED_TEXT_INDEX = "partial_nested_text"; - private static final String NESTED_PATTERN = "partial_nested_*"; - - // Truncation fixture: 1 keyword index + 8 bare-text indices, so the excluded list exceeds the - // warning's spell-out cap and must be summarized as "... and N more". - private static final String MANY_KEYWORD_INDEX = "partial_many_keyword"; - private static final String MANY_TEXT_PREFIX = "partial_many_text"; - private static final String MANY_PATTERN = "partial_many_*"; - private static final int MANY_TEXT_COUNT = 8; - - // Priority-ladder fixture: one keyword index vs two text-with-.keyword indices. Keyword is - // outnumbered, so a count-based majority would keep the text-with-.keyword group; the - // deterministic keyword-first rule must keep the single keyword index instead. - private static final String PRIORITY_KEYWORD_INDEX = "partial_priority_keyword"; - private static final String PRIORITY_TEXTKW_INDEX_1 = "partial_priority_textkw1"; - private static final String PRIORITY_TEXTKW_INDEX_2 = "partial_priority_textkw2"; - private static final String PRIORITY_PATTERN = "partial_priority_*"; - - // Multi-field expression fixture: two fields, both keyword in one index and both bare text in - // another. A group key like concat(city, region) must trace to BOTH fields and keep only the - // index where both are aggregatable. - private static final String MULTI_KEYWORD_INDEX = "partial_multi_keyword"; - private static final String MULTI_TEXT_INDEX = "partial_multi_text"; - private static final String MULTI_PATTERN = "partial_multi_*"; - - // Non-text-type fixture: the field is an aggregatable integer in one index and bare text in - // another. The integer index is kept and the text index excluded, rather than the field silently - // coercing to one type and dropping the other index's docs. - private static final String NUMTEXT_INT_INDEX = "partial_numtext_int"; - private static final String NUMTEXT_TEXT_INDEX = "partial_numtext_text"; - private static final String NUMTEXT_PATTERN = "partial_numtext_*"; - - @Override - public void init() throws Exception { - super.init(); - enableCalcite(); - createTestIndices(); - } - - @After - public void cleanup() throws IOException { - setPartialResult(false); - setPitContextLimit(null); - } - - private void createTestIndices() throws IOException { - // keyword index: env is aggregatable. Two shards so a scan needs 2 PIT contexts. - if (!isIndexExist(client(), KEYWORD_INDEX)) { - String mapping = - "{\"settings\":{\"index\":{\"number_of_shards\":2,\"number_of_replicas\":0}}," - + "\"mappings\":{\"properties\":{\"env\":{\"type\":\"keyword\"}}}}"; - createIndexByRestClient(client(), KEYWORD_INDEX, mapping); - Request bulk = new Request("POST", "/" + KEYWORD_INDEX + "/_bulk?refresh=true"); - bulk.setJsonEntity( - "{\"index\":{}}\n{\"env\":\"prod\"}\n" - + "{\"index\":{}}\n{\"env\":\"prod\"}\n" - + "{\"index\":{}}\n{\"env\":\"dev\"}\n"); - performRequest(client(), bulk); - } - // text index (no .keyword sub-field): env is NOT aggregatable -> forces the conflict collapse. - if (!isIndexExist(client(), TEXT_INDEX)) { - String mapping = - "{\"settings\":{\"index\":{\"number_of_shards\":2,\"number_of_replicas\":0}}," - + "\"mappings\":{\"properties\":{\"env\":{\"type\":\"text\"}}}}"; - createIndexByRestClient(client(), TEXT_INDEX, mapping); - Request bulk = new Request("POST", "/" + TEXT_INDEX + "/_bulk?refresh=true"); - bulk.setJsonEntity( - "{\"index\":{}}\n{\"env\":\"prod\"}\n" + "{\"index\":{}}\n{\"env\":\"qa\"}\n"); - performRequest(client(), bulk); - } - - // A nested/dotted field (resource.attributes.env) is stored as an object tree in the mapping, - // so the partitioning must flatten it to match the bucket field's dotted path. Mirrors the - // real observability shape (e.g. resource.attributes.applicationid). - if (!isIndexExist(client(), NESTED_KEYWORD_INDEX)) { - String mapping = - "{\"settings\":{\"index\":{\"number_of_shards\":2,\"number_of_replicas\":0}}," - + "\"mappings\":{\"properties\":{\"resource\":{\"properties\":{\"attributes\":" - + "{\"properties\":{\"env\":{\"type\":\"keyword\"}}}}}}}}"; - createIndexByRestClient(client(), NESTED_KEYWORD_INDEX, mapping); - Request bulk = new Request("POST", "/" + NESTED_KEYWORD_INDEX + "/_bulk?refresh=true"); - bulk.setJsonEntity( - "{\"index\":{}}\n{\"resource\":{\"attributes\":{\"env\":\"prod\"}}}\n" - + "{\"index\":{}}\n{\"resource\":{\"attributes\":{\"env\":\"prod\"}}}\n" - + "{\"index\":{}}\n{\"resource\":{\"attributes\":{\"env\":\"dev\"}}}\n"); - performRequest(client(), bulk); - } - if (!isIndexExist(client(), NESTED_TEXT_INDEX)) { - String mapping = - "{\"settings\":{\"index\":{\"number_of_shards\":2,\"number_of_replicas\":0}}," - + "\"mappings\":{\"properties\":{\"resource\":{\"properties\":{\"attributes\":" - + "{\"properties\":{\"env\":{\"type\":\"text\"}}}}}}}}"; - createIndexByRestClient(client(), NESTED_TEXT_INDEX, mapping); - Request bulk = new Request("POST", "/" + NESTED_TEXT_INDEX + "/_bulk?refresh=true"); - bulk.setJsonEntity( - "{\"index\":{}}\n{\"resource\":{\"attributes\":{\"env\":\"prod\"}}}\n" - + "{\"index\":{}}\n{\"resource\":{\"attributes\":{\"env\":\"qa\"}}}\n"); - performRequest(client(), bulk); - } - - // Priority ladder: 1 keyword index vs 2 text-with-.keyword indices (keyword outnumbered 2:1). - if (!isIndexExist(client(), PRIORITY_KEYWORD_INDEX)) { - String mapping = - "{\"settings\":{\"index\":{\"number_of_shards\":2,\"number_of_replicas\":0}}," - + "\"mappings\":{\"properties\":{\"env\":{\"type\":\"keyword\"}}}}"; - createIndexByRestClient(client(), PRIORITY_KEYWORD_INDEX, mapping); - Request bulk = new Request("POST", "/" + PRIORITY_KEYWORD_INDEX + "/_bulk?refresh=true"); - bulk.setJsonEntity( - "{\"index\":{}}\n{\"env\":\"prod\"}\n" + "{\"index\":{}}\n{\"env\":\"dev\"}\n"); - performRequest(client(), bulk); - } - String textKwMapping = - "{\"settings\":{\"index\":{\"number_of_shards\":2,\"number_of_replicas\":0}}," - + "\"mappings\":{\"properties\":{\"env\":{\"type\":\"text\",\"fields\":" - + "{\"keyword\":{\"type\":\"keyword\",\"ignore_above\":256}}}}}}"; - for (String idx : new String[] {PRIORITY_TEXTKW_INDEX_1, PRIORITY_TEXTKW_INDEX_2}) { - if (!isIndexExist(client(), idx)) { - createIndexByRestClient(client(), idx, textKwMapping); - Request bulk = new Request("POST", "/" + idx + "/_bulk?refresh=true"); - bulk.setJsonEntity( - "{\"index\":{}}\n{\"env\":\"prod\"}\n" + "{\"index\":{}}\n{\"env\":\"stage\"}\n"); - performRequest(client(), bulk); - } - } - - // Truncation: 1 keyword + many bare-text indices, so the excluded list exceeds the warning cap. - if (!isIndexExist(client(), MANY_KEYWORD_INDEX)) { - String mapping = - "{\"settings\":{\"index\":{\"number_of_shards\":2,\"number_of_replicas\":0}}," - + "\"mappings\":{\"properties\":{\"env\":{\"type\":\"keyword\"}}}}"; - createIndexByRestClient(client(), MANY_KEYWORD_INDEX, mapping); - Request bulk = new Request("POST", "/" + MANY_KEYWORD_INDEX + "/_bulk?refresh=true"); - bulk.setJsonEntity("{\"index\":{}}\n{\"env\":\"prod\"}\n"); - performRequest(client(), bulk); - } - // Multi-field expression fixture: both fields keyword in one index, both bare text in another. - if (!isIndexExist(client(), MULTI_KEYWORD_INDEX)) { - String mapping = - "{\"settings\":{\"index\":{\"number_of_shards\":2,\"number_of_replicas\":0}}," - + "\"mappings\":{\"properties\":{\"city\":{\"type\":\"keyword\"}," - + "\"region\":{\"type\":\"keyword\"}}}}"; - createIndexByRestClient(client(), MULTI_KEYWORD_INDEX, mapping); - Request bulk = new Request("POST", "/" + MULTI_KEYWORD_INDEX + "/_bulk?refresh=true"); - bulk.setJsonEntity( - "{\"index\":{}}\n{\"city\":\"nyc\",\"region\":\"us\"}\n" - + "{\"index\":{}}\n{\"city\":\"nyc\",\"region\":\"us\"}\n" - + "{\"index\":{}}\n{\"city\":\"sf\",\"region\":\"us\"}\n"); - performRequest(client(), bulk); - } - if (!isIndexExist(client(), MULTI_TEXT_INDEX)) { - String mapping = - "{\"settings\":{\"index\":{\"number_of_shards\":2,\"number_of_replicas\":0}}," - + "\"mappings\":{\"properties\":{\"city\":{\"type\":\"text\"}," - + "\"region\":{\"type\":\"text\"}}}}"; - createIndexByRestClient(client(), MULTI_TEXT_INDEX, mapping); - Request bulk = new Request("POST", "/" + MULTI_TEXT_INDEX + "/_bulk?refresh=true"); - bulk.setJsonEntity( - "{\"index\":{}}\n{\"city\":\"la\",\"region\":\"us\"}\n" - + "{\"index\":{}}\n{\"city\":\"sea\",\"region\":\"us\"}\n"); - performRequest(client(), bulk); - } - - // Non-text-type conflict: integer vs bare text on the same field. - if (!isIndexExist(client(), NUMTEXT_INT_INDEX)) { - String mapping = - "{\"settings\":{\"index\":{\"number_of_shards\":2,\"number_of_replicas\":0}}," - + "\"mappings\":{\"properties\":{\"val\":{\"type\":\"integer\"}}}}"; - createIndexByRestClient(client(), NUMTEXT_INT_INDEX, mapping); - Request bulk = new Request("POST", "/" + NUMTEXT_INT_INDEX + "/_bulk?refresh=true"); - bulk.setJsonEntity( - "{\"index\":{}}\n{\"val\":7}\n" - + "{\"index\":{}}\n{\"val\":7}\n" - + "{\"index\":{}}\n{\"val\":9}\n"); - performRequest(client(), bulk); - } - if (!isIndexExist(client(), NUMTEXT_TEXT_INDEX)) { - String mapping = - "{\"settings\":{\"index\":{\"number_of_shards\":2,\"number_of_replicas\":0}}," - + "\"mappings\":{\"properties\":{\"val\":{\"type\":\"text\"}}}}"; - createIndexByRestClient(client(), NUMTEXT_TEXT_INDEX, mapping); - Request bulk = new Request("POST", "/" + NUMTEXT_TEXT_INDEX + "/_bulk?refresh=true"); - bulk.setJsonEntity( - "{\"index\":{}}\n{\"val\":\"aa\"}\n" + "{\"index\":{}}\n{\"val\":\"bb\"}\n"); - performRequest(client(), bulk); - } - - String bareTextMapping = - "{\"settings\":{\"index\":{\"number_of_shards\":2,\"number_of_replicas\":0}}," - + "\"mappings\":{\"properties\":{\"env\":{\"type\":\"text\"}}}}"; - for (int i = 1; i <= MANY_TEXT_COUNT; i++) { - String idx = MANY_TEXT_PREFIX + i; - if (!isIndexExist(client(), idx)) { - createIndexByRestClient(client(), idx, bareTextMapping); - Request bulk = new Request("POST", "/" + idx + "/_bulk?refresh=true"); - bulk.setJsonEntity("{\"index\":{}}\n{\"env\":\"prod\"}\n"); - performRequest(client(), bulk); - } - } - } - - @Test - public void partialResultOffReturnsCompleteResultWithoutWarning() throws IOException { - setPartialResult(false); - // Since #5646 the collapsed text group key pushes down as a per-document _source script, so the - // complete answer is returned (slowly) rather than failing. Every index contributes: the - // keyword - // index (prod=2, dev=1) plus the text index (prod=1, qa=1). - JSONObject result = - executeQuery(String.format("source=%s | stats count() by env | sort env", PATTERN)); - verifyDataRows(result, rows(1, "dev"), rows(1, "qa"), rows(3, "prod")); - assertTrue("a complete result carries no partial-result warning", !result.has("warnings")); - } - - @Test - @RequiresCapability(CROSS_INDEX_INCOMPATIBLE_TYPES) - public void partialResultOnReturnsKeywordSubsetWithWarning() throws IOException { - setPartialResult(true); - // Even with the PIT budget crippled, partial mode pushes the aggregation down (size=0), so no - // PIT is opened and the query succeeds over the aggregatable keyword index only. - setPitContextLimit("1"); - JSONObject result = - executeQuery(String.format("source=%s | stats count() by env | sort env", PATTERN)); - - // Only the keyword index contributes: prod=2, dev=1. The text index (prod=1, qa=1) is excluded. - verifyDataRows(result, rows(1, "dev"), rows(2, "prod")); - - assertTrue("response should carry a warnings array", result.has("warnings")); - JSONArray warnings = result.getJSONArray("warnings"); - assertEquals(1, warnings.length()); - JSONObject warning = warnings.getJSONObject(0); - assertEquals("PARTIAL_RESULT", warning.getString("type")); - assertTrue( - "warning detail should name the excluded text index", - warning.getString("detail").contains(TEXT_INDEX)); - } - - @Test - public void partialResultOffOmitsWarningWhenNoConflict() throws IOException { - setPartialResult(true); - setPitContextLimit(null); - // A single-index aggregatable query has no conflict, so it pushes down normally and no warning - // is attached even with partial mode enabled. - JSONObject result = - executeQuery(String.format("source=%s | stats count() by env | sort env", KEYWORD_INDEX)); - verifyDataRows(result, rows(1, "dev"), rows(2, "prod")); - assertTrue("no warning expected on a clean aggregation", !result.has("warnings")); - } - - @Test - @RequiresCapability(CROSS_INDEX_INCOMPATIBLE_TYPES) - public void partialResultOnHandlesNestedDottedField() throws IOException { - setPartialResult(true); - setPitContextLimit("1"); - // The grouped field is a nested/dotted path; the partitioning must flatten the mapping to find - // it. Only the keyword index contributes: prod=2, dev=1. - JSONObject result = - executeQuery( - String.format( - "source=%s | stats count() by resource.attributes.env | sort" - + " `resource.attributes.env`", - NESTED_PATTERN)); - verifyDataRows(result, rows(1, "dev"), rows(2, "prod")); - - assertTrue("response should carry a warnings array", result.has("warnings")); - JSONObject warning = result.getJSONArray("warnings").getJSONObject(0); - assertEquals("PARTIAL_RESULT", warning.getString("type")); - assertTrue( - "warning should name the dotted field", - warning.getString("detail").contains("resource.attributes.env")); - assertTrue( - "warning should name the excluded nested-text index", - warning.getString("detail").contains(NESTED_TEXT_INDEX)); - } - - @Test - @RequiresCapability(CROSS_INDEX_INCOMPATIBLE_TYPES) - public void partialResultOnHandlesEvalDerivedGroupKey() throws IOException { - setPartialResult(true); - setPitContextLimit("1"); - // The group key is an expression over the conflicting field (upper(env)), not the bare field. - // Partitioning traces it back to env, so the keyword index is kept and the text index excluded - // just as for a bare group key. Only the keyword index contributes: PROD=2, DEV=1. - JSONObject result = - executeQuery( - String.format( - "source=%s | eval g = upper(env) | stats count() by g | sort g", PATTERN)); - verifyDataRows(result, rows(1, "DEV"), rows(2, "PROD")); - - assertTrue("response should carry a warnings array", result.has("warnings")); - JSONObject warning = result.getJSONArray("warnings").getJSONObject(0); - assertEquals("PARTIAL_RESULT", warning.getString("type")); - assertTrue( - "warning should name the underlying field the expression reads", - warning.getString("detail").contains("env")); - assertTrue( - "warning should name the excluded text index", - warning.getString("detail").contains(TEXT_INDEX)); - } - - @Test - @RequiresCapability(CROSS_INDEX_INCOMPATIBLE_TYPES) - public void partialResultOnHandlesMultiFieldExpressionGroupKey() throws IOException { - setPartialResult(true); - setPitContextLimit("1"); - // The group key reads two fields (concat(city, region)); partitioning must trace it to BOTH and - // keep only the index where both are aggregatable. Keyword index: nycus=2, sfus=1; text - // excluded. - JSONObject result = - executeQuery( - String.format( - "source=%s | eval g = concat(city, region) | stats count() by g | sort g", - MULTI_PATTERN)); - verifyDataRows(result, rows(2, "nycus"), rows(1, "sfus")); - - JSONObject warning = result.getJSONArray("warnings").getJSONObject(0); - assertEquals("PARTIAL_RESULT", warning.getString("type")); - assertTrue( - "warning should name both underlying fields the expression reads", - warning.getString("detail").contains("city") - && warning.getString("detail").contains("region")); - assertTrue( - "warning should name the excluded text index", - warning.getString("detail").contains(MULTI_TEXT_INDEX)); - } - - @Test - @RequiresCapability(CROSS_INDEX_INCOMPATIBLE_TYPES) - public void partialResultKeepsNumericExcludesText() throws IOException { - setPartialResult(true); - setPitContextLimit("1"); - // The field is integer in one index and bare text in another. The integer index is aggregatable - // so it is kept and the text index excluded (with a warning), instead of silently dropping one - // index's docs to a coerced type. The bucket labels materialize correctly (no null). - JSONObject result = - executeQuery( - String.format("source=%s | stats count() as c by val | sort val", NUMTEXT_PATTERN)); - // Only the integer index contributes its two buckets; the text values (aa, bb) are excluded. - assertEquals(2, result.getJSONArray("datarows").length()); - String body = result.toString(); - assertTrue("excluded text values must not appear: " + body, !body.contains("aa")); - assertTrue("excluded text values must not appear: " + body, !body.contains("bb")); - - JSONObject warning = result.getJSONArray("warnings").getJSONObject(0); - assertEquals("PARTIAL_RESULT", warning.getString("type")); - assertTrue( - "warning should name the excluded text index", - warning.getString("detail").contains(NUMTEXT_TEXT_INDEX)); - } - - @Test - @RequiresCapability(CROSS_INDEX_INCOMPATIBLE_TYPES) - public void partialResultKeepsKeywordGroupEvenWhenOutnumbered() throws IOException { - setPartialResult(true); - setPitContextLimit("1"); - // Keyword is outnumbered 2:1 by text-with-.keyword indices. The deterministic keyword-first - // rule keeps the single keyword index (prod:1, dev:1) and excludes both text-with-.keyword - // indices -- a count-based majority would have kept the text group instead. - JSONObject result = - executeQuery( - String.format("source=%s | stats count() by env | sort env", PRIORITY_PATTERN)); - verifyDataRows(result, rows(1, "dev"), rows(1, "prod")); - - JSONObject warning = result.getJSONArray("warnings").getJSONObject(0); - assertEquals("PARTIAL_RESULT", warning.getString("type")); - assertTrue( - "both text-with-keyword indices should be excluded", - warning.getString("detail").contains(PRIORITY_TEXTKW_INDEX_1) - && warning.getString("detail").contains(PRIORITY_TEXTKW_INDEX_2)); - } - - @Test - @RequiresCapability(CROSS_INDEX_INCOMPATIBLE_TYPES) - public void partialResultWarningTruncatesLargeExcludedList() throws IOException { - setPartialResult(true); - setPitContextLimit("1"); - JSONObject result = - executeQuery(String.format("source=%s | stats count() by env", MANY_PATTERN)); - - JSONObject warning = result.getJSONArray("warnings").getJSONObject(0); - // 8 bare-text indices excluded; the message reports the exact count... - assertTrue( - "message should report the full excluded count", - warning.getString("message").contains("8 of 9")); - // ...but the detail spells out only a few and summarizes the rest. - String detail = warning.getString("detail"); - assertTrue("detail should summarize the remainder", detail.contains("and 3 more")); - assertTrue( - "detail should not list every excluded index", !detail.contains(MANY_TEXT_PREFIX + "8")); - } - - @Test - @RequiresCapability(CROSS_INDEX_INCOMPATIBLE_TYPES) - public void partialResultRefusedForCsvFormat() throws IOException { - setPartialResult(true); - // CSV has no warnings channel, so partial mode must NOT silently drop the text index -- there - // would be no way to tell the caller the numbers are undercounted. It falls through to the - // normal (complete) path instead, so the text index's rows are still counted. - String csv = - executeCsvQuery( - String.format("source=%s | stats count() by env | sort env", PATTERN), false); - // qa exists only in the excluded text index: its presence proves nothing was dropped. - assertTrue( - "CSV must return the complete result, including the text index: " + csv, - csv.contains("qa")); - } - - private void setPartialResult(boolean enabled) throws IOException { - updateClusterSettings( - new ClusterSetting( - "persistent", - Settings.Key.PARTIAL_RESULT_ON_MAPPING_CONFLICT.getKeyValue(), - Boolean.toString(enabled))); - } - - private void setPitContextLimit(String value) throws IOException { - updateClusterSettings(new ClusterSetting("transient", "search.max_open_pit_context", value)); - } -} diff --git a/integ-test/src/test/java/org/opensearch/sql/security/PartialResultSecurityIT.java b/integ-test/src/test/java/org/opensearch/sql/security/PartialResultSecurityIT.java deleted file mode 100644 index c8296387183..00000000000 --- a/integ-test/src/test/java/org/opensearch/sql/security/PartialResultSecurityIT.java +++ /dev/null @@ -1,169 +0,0 @@ -/* - * Copyright OpenSearch Contributors - * SPDX-License-Identifier: Apache-2.0 - */ - -package org.opensearch.sql.security; - -import static org.opensearch.sql.util.MatcherUtils.rows; -import static org.opensearch.sql.util.MatcherUtils.verifyDataRows; -import static org.opensearch.sql.util.TestUtils.createIndexByRestClient; -import static org.opensearch.sql.util.TestUtils.isIndexExist; -import static org.opensearch.sql.util.TestUtils.performRequest; - -import java.io.IOException; -import java.util.Locale; -import org.json.JSONArray; -import org.json.JSONObject; -import org.junit.After; -import org.junit.Test; -import org.opensearch.client.Request; -import org.opensearch.client.RequestOptions; -import org.opensearch.client.Response; -import org.opensearch.sql.common.setting.Settings; -import org.opensearch.sql.util.ClusterPlugins; - -/** - * Runs the partial-result-on-mapping-conflict path with the security plugin installed. Regression - * guard for #5739: the warnings-supported gate used to live in Log4j {@code ThreadContext}, which - * the security plugin's transport interceptor drops on the transport-to-worker handoff, so partial - * mode silently bailed and returned a complete result with no warning. The {@code - * integTestWithSecurity} suite never exercised this path -- it only runs {@code - * org.opensearch.sql.security.*}, and the partial-result IT lives elsewhere -- so only the release - * distribution's full-suite-with-security caught it. Placing this test in the security package - * closes that gap. - */ -public class PartialResultSecurityIT extends SecurityTestBase { - - private static final String KEYWORD_INDEX = "partial_sec_keyword"; - private static final String TEXT_INDEX = "partial_sec_text"; - private static final String PATTERN = "partial_sec_*"; - - private static final String USER = "partial_sec_user"; - private static final String ROLE = "partial_sec_role"; - - private boolean initialized = false; - - @Override - protected void init() throws Exception { - ClusterPlugins.requirePluginOrAssume( - client(), - ClusterPlugins.SECURITY_PLUGIN, - "opensearch-security plugin not installed on test cluster; skipping FGAC tests"); - super.init(); - enableCalcite(); - if (!initialized) { - createRoleWithIndexAccess(ROLE, PATTERN); - createUser(USER, ROLE); - createConflictIndices(); - initialized = true; - } - } - - @After - public void resetPartialResult() throws IOException { - setPartialResult(false); - } - - private void createConflictIndices() throws IOException { - // env is an aggregatable keyword here... - if (!isIndexExist(client(), KEYWORD_INDEX)) { - String mapping = "{\"mappings\":{\"properties\":{\"env\":{\"type\":\"keyword\"}}}}"; - createIndexByRestClient(client(), KEYWORD_INDEX, mapping); - Request bulk = new Request("POST", "/" + KEYWORD_INDEX + "/_bulk?refresh=true"); - bulk.setJsonEntity( - "{\"index\":{}}\n{\"env\":\"prod\"}\n" - + "{\"index\":{}}\n{\"env\":\"prod\"}\n" - + "{\"index\":{}}\n{\"env\":\"dev\"}\n"); - performRequest(client(), bulk); - } - // ...and bare text (no .keyword sub-field) here, so the field collapses to non-aggregatable. - if (!isIndexExist(client(), TEXT_INDEX)) { - String mapping = "{\"mappings\":{\"properties\":{\"env\":{\"type\":\"text\"}}}}"; - createIndexByRestClient(client(), TEXT_INDEX, mapping); - Request bulk = new Request("POST", "/" + TEXT_INDEX + "/_bulk?refresh=true"); - bulk.setJsonEntity( - "{\"index\":{}}\n{\"env\":\"prod\"}\n" + "{\"index\":{}}\n{\"env\":\"qa\"}\n"); - performRequest(client(), bulk); - } - } - - @Test - public void partialResultWarningSurvivesSecurityHandoff() throws IOException { - setPartialResult(true); - JSONObject result = - executeQueryAsUser( - String.format("source=%s | stats count() by env | sort env", PATTERN), USER); - - // Only the aggregatable keyword index contributes (prod=2, dev=1); the text index is excluded. - verifyDataRows(result, rows(1, "dev"), rows(2, "prod")); - - assertTrue( - "partial result must carry a warnings channel under security", result.has("warnings")); - JSONArray warnings = result.getJSONArray("warnings"); - assertEquals(1, warnings.length()); - JSONObject warning = warnings.getJSONObject(0); - assertEquals("PARTIAL_RESULT", warning.getString("type")); - assertTrue( - "warning should name the excluded text index", - warning.getString("detail").contains(TEXT_INDEX)); - } - - @Test - public void completeResultCarriesNoWarningWithSecurity() throws IOException { - setPartialResult(false); - JSONObject result = - executeQueryAsUser( - String.format("source=%s | stats count() by env | sort env", PATTERN), USER); - // Every index contributes: keyword (prod=2, dev=1) + text (prod=1, qa=1). - verifyDataRows(result, rows(1, "dev"), rows(1, "qa"), rows(3, "prod")); - assertFalse("a complete result carries no warning", result.has("warnings")); - } - - @Test - public void perRequestPartialResultFalseOverridesClusterSettingUnderSecurity() - throws IOException { - // Cluster setting ON, but the request explicitly opts OUT via partial_result=false. The - // per-request override must win -> complete result over all indices, no warning. On the buggy - // code the override lives in Log4j ThreadContext (QueryContext.setPartialResultOverride) and is - // dropped by the security transport->worker handoff -- the same drop #5739 fixed for - // warningsSupported but left in place for the override -- so it silently falls back to the ON - // cluster setting and returns a partial result with a warning. - setPartialResult(true); - JSONObject result = - executeQueryAsUserWithPartialResult( - String.format("source=%s | stats count() by env | sort env", PATTERN), USER, false); - // Every index contributes: keyword (prod=2, dev=1) + text (prod=1, qa=1). - verifyDataRows(result, rows(1, "dev"), rows(1, "qa"), rows(3, "prod")); - assertFalse( - "partial_result=false must override the ON cluster setting -> complete result, no warning", - result.has("warnings")); - } - - /** Like {@link #executeQueryAsUser}, but also sends the per-request {@code partial_result}. */ - private JSONObject executeQueryAsUserWithPartialResult( - String query, String username, boolean partialResult) throws IOException { - Request request = new Request("POST", "/_plugins/_ppl"); - request.setJsonEntity( - String.format( - Locale.ROOT, - "{ \"query\": \"%s\", \"partial_result\": %s }", - query, - Boolean.toString(partialResult))); - RequestOptions.Builder options = RequestOptions.DEFAULT.toBuilder(); - options.addHeader("Content-Type", "application/json"); - options.addHeader("Authorization", createBasicAuthHeader(username, STRONG_PASSWORD)); - request.setOptions(options); - Response response = client().performRequest(request); - assertEquals(200, response.getStatusLine().getStatusCode()); - return new JSONObject(org.opensearch.sql.legacy.TestUtils.getResponseBody(response, true)); - } - - private void setPartialResult(boolean enabled) throws IOException { - updateClusterSettings( - new ClusterSetting( - "persistent", - Settings.Key.PARTIAL_RESULT_ON_MAPPING_CONFLICT.getKeyValue(), - Boolean.toString(enabled))); - } -} diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/data/value/OpenSearchExprValueFactory.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/data/value/OpenSearchExprValueFactory.java index 9822b6f411d..aeda2da0ce6 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/data/value/OpenSearchExprValueFactory.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/data/value/OpenSearchExprValueFactory.java @@ -247,17 +247,15 @@ private ExprValue parse( } /** - * String form of scalar content destined for a text/keyword column. A numeric or boolean bucket - * key can land in a string column -- e.g. a partial-result aggregation over an index where the - * field is numeric while the conflict's merged type is text. Render it as its string form rather - * than letting the {@code (String) value} cast fail and null the value out. + * String form of scalar content destined for a text/keyword column. A numeric or boolean value + * can land in a string column -- e.g. an aggregation over indices where the field is numeric in + * one and keyword in another, so some bucket keys come back as numbers. Render it as its string + * form rather than letting the {@code (String) value} cast fail and null the value out. */ private static String stringOf(Content content) { try { return content.stringValue(); } catch (RuntimeException e) { - // Not a string value (e.g. a numeric aggregation bucket key landing in a text column via a - // partial-result narrowing) -- render its string form instead of failing the cast to null. return String.valueOf(content.objectValue()); } } diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/request/system/OpenSearchDescribeIndexRequest.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/request/system/OpenSearchDescribeIndexRequest.java index b66a6281f60..2b65b5c7af1 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/request/system/OpenSearchDescribeIndexRequest.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/request/system/OpenSearchDescribeIndexRequest.java @@ -15,7 +15,6 @@ import java.util.List; import java.util.Locale; import java.util.Map; -import lombok.Getter; import lombok.extern.log4j.Log4j2; import org.opensearch.sql.data.model.ExprTupleValue; import org.opensearch.sql.data.model.ExprValue; @@ -93,14 +92,6 @@ public List search() { return results; } - /** - * The per-index mappings behind the last {@link #getFieldTypes()} call, keyed by concrete index - * name. Retained because merging discards which index mapped a field which way, and callers that - * need that detail (e.g. partitioning a wildcard by whether a field is aggregatable) would - * otherwise have to fetch the mappings a second time. - */ - @Getter private Map lastIndexMappings = Map.of(); - /** * Get the mapping of field and type. * @@ -111,14 +102,13 @@ public Map getFieldTypes() { Map fieldTypes = new HashMap<>(); Map indexMappings = client.getIndexMappings(getLocalIndexNames(indexName.getIndexNames())); - this.lastIndexMappings = indexMappings; if (indexMappings.size() <= 1) { for (IndexMapping indexMapping : indexMappings.values()) { fieldTypes.putAll(indexMapping.getFieldMappings()); } } else { - // Merge deep copies: MergeRuleHelper mutates the field mappings in place, and the per-index - // mappings retained above must stay intact for partial-result partitioning. + // Merge deep copies: MergeRuleHelper mutates the field mappings in place, and a fetched + // mapping must not be altered by being merged. for (IndexMapping indexMapping : indexMappings.values()) { MergeRuleHelper.merge(fieldTypes, deepCopy(indexMapping.getFieldMappings())); } diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/setting/OpenSearchSettings.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/setting/OpenSearchSettings.java index ad6bdcb7c9c..e7e4232bb1a 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/setting/OpenSearchSettings.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/setting/OpenSearchSettings.java @@ -195,13 +195,6 @@ public class OpenSearchSettings extends Settings { Setting.Property.NodeScope, Setting.Property.Dynamic); - public static final Setting PARTIAL_RESULT_ON_MAPPING_CONFLICT_SETTING = - Setting.boolSetting( - Key.PARTIAL_RESULT_ON_MAPPING_CONFLICT.getKeyValue(), - false, - Setting.Property.NodeScope, - Setting.Property.Dynamic); - public static final Setting QUERY_MEMORY_LIMIT_SETTING = Setting.memorySizeSetting( Key.QUERY_MEMORY_LIMIT.getKeyValue(), @@ -534,12 +527,6 @@ public OpenSearchSettings(ClusterSettings clusterSettings) { Key.CALCITE_SUPPORT_ALL_JOIN_TYPES, CALCITE_SUPPORT_ALL_JOIN_TYPES_SETTING, new Updater(Key.CALCITE_SUPPORT_ALL_JOIN_TYPES)); - register( - settingBuilder, - clusterSettings, - Key.PARTIAL_RESULT_ON_MAPPING_CONFLICT, - PARTIAL_RESULT_ON_MAPPING_CONFLICT_SETTING, - new Updater(Key.PARTIAL_RESULT_ON_MAPPING_CONFLICT)); register( settingBuilder, clusterSettings, @@ -773,7 +760,6 @@ public static List> pluginSettings() { .add(CALCITE_PUSHDOWN_ENABLED_SETTING) .add(CALCITE_PUSHDOWN_ROWCOUNT_ESTIMATION_FACTOR_SETTING) .add(CALCITE_SUPPORT_ALL_JOIN_TYPES_SETTING) - .add(PARTIAL_RESULT_ON_MAPPING_CONFLICT_SETTING) .add(DEFAULT_PATTERN_METHOD_SETTING) .add(DEFAULT_PATTERN_MODE_SETTING) .add(DEFAULT_PATTERN_MAX_SAMPLE_COUNT_SETTING) diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/OpenSearchIndex.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/OpenSearchIndex.java index 8a56fc24e15..e90eeace6e3 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/OpenSearchIndex.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/OpenSearchIndex.java @@ -29,7 +29,6 @@ import org.opensearch.sql.opensearch.client.OpenSearchClient; import org.opensearch.sql.opensearch.data.type.OpenSearchDataType; import org.opensearch.sql.opensearch.data.value.OpenSearchExprValueFactory; -import org.opensearch.sql.opensearch.mapping.IndexMapping; import org.opensearch.sql.opensearch.monitor.OpenSearchMemoryHealthy; import org.opensearch.sql.opensearch.monitor.OpenSearchResourceMonitor; import org.opensearch.sql.opensearch.planner.physical.ADOperator; @@ -92,13 +91,6 @@ public class OpenSearchIndex extends AbstractOpenSearchTable { /** The cached mapping of alias type field to its original path. */ private Map aliasMapping = null; - /** - * The cached per-index field mappings, keyed by concrete index name. Populated as a by-product of - * resolving {@link #cachedFieldOpenSearchTypes}, since merging those types discards which index - * mapped a field which way. - */ - private Map cachedIndexMappings = null; - /** The cached max result window setting of index. */ private Integer cachedMaxResultWindow = null; @@ -190,24 +182,12 @@ public Map getFieldOpenSearchTypes() { return cachedFieldOpenSearchTypes; } - /** - * The per-index field mappings behind this index's merged types, keyed by concrete index name - * (the wildcard, if any, is already resolved). Needed by callers that must know which index - * mapped a field which way -- the merged view in {@link #getFieldOpenSearchTypes()} discards - * that. Shares the mapping fetch with the merged types, so this costs no extra round trip. - */ - public Map getIndexMappings() { - resolveFieldOpenSearchTypes(); - return cachedIndexMappings; - } - - /** Fetch and cache the merged field types, retaining the per-index mappings behind them. */ + /** Fetch and cache the merged field types. */ private void resolveFieldOpenSearchTypes() { if (cachedFieldOpenSearchTypes == null) { OpenSearchDescribeIndexRequest request = new OpenSearchDescribeIndexRequest(client, indexName); cachedFieldOpenSearchTypes = request.getFieldTypes(); - cachedIndexMappings = request.getLastIndexMappings(); } } diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteLogicalIndexScan.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteLogicalIndexScan.java index 09cd5278dc5..2017437e7bd 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteLogicalIndexScan.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteLogicalIndexScan.java @@ -7,11 +7,9 @@ import com.google.common.collect.ImmutableList; import java.util.ArrayList; -import java.util.LinkedHashSet; import java.util.List; import java.util.Map; import java.util.Objects; -import java.util.Set; import java.util.stream.Collectors; import javax.annotation.Nullable; import lombok.Getter; @@ -36,9 +34,7 @@ import org.apache.calcite.rel.type.RelDataTypeFactory; import org.apache.calcite.rel.type.RelDataTypeField; import org.apache.calcite.rex.RexBuilder; -import org.apache.calcite.rex.RexInputRef; import org.apache.calcite.rex.RexNode; -import org.apache.calcite.rex.RexVisitorImpl; import org.apache.calcite.sql.fun.SqlStdOperatorTable; import org.apache.calcite.sql.type.SqlTypeName; import org.apache.commons.lang3.tuple.Pair; @@ -46,7 +42,6 @@ import org.apache.logging.log4j.Logger; import org.opensearch.search.aggregations.AggregationBuilder; import org.opensearch.sql.ast.tree.HighlightConfig; -import org.opensearch.sql.calcite.CalcitePlanContext; import org.opensearch.sql.calcite.plan.HighlightPushDown; import org.opensearch.sql.calcite.utils.OpenSearchTypeFactory; import org.opensearch.sql.calcite.utils.PPLHintUtils; @@ -56,7 +51,6 @@ import org.opensearch.sql.expression.HighlightExpression; import org.opensearch.sql.opensearch.data.type.OpenSearchDataType; import org.opensearch.sql.opensearch.data.type.OpenSearchTextType; -import org.opensearch.sql.opensearch.mapping.IndexMapping; import org.opensearch.sql.opensearch.planner.rules.OpenSearchIndexRules; import org.opensearch.sql.opensearch.request.AggregateAnalyzer; import org.opensearch.sql.opensearch.request.PredicateAnalyzer; @@ -378,15 +372,6 @@ public CalciteLogicalIndexScan pushDownRareTop(Project project, RareTopDigest di } public AbstractRelNode pushDownAggregate(Aggregate aggregate, @Nullable Project project) { - return pushDownAggregate(aggregate, project, true); - } - - /** - * @param allowPartialFallback whether the partial-result path may be tried. False when - * re-entering from that path over the narrowed index subset, so it is attempted at most once. - */ - private AbstractRelNode pushDownAggregate( - Aggregate aggregate, @Nullable Project project, boolean allowPartialFallback) { try { CalciteLogicalIndexScan newScan = new CalciteLogicalIndexScan( @@ -414,17 +399,6 @@ private AbstractRelNode pushDownAggregate( } return null; } - // Try partial mode before analyze: since #5646 a text/keyword conflict pushes down as a slow - // _source script instead of failing, so a post-failure fallback would never fire. - if (allowPartialFallback) { - List partitionFields = resolvePartitionFields(aggregate, project); - if (partitionFields != null) { - AbstractRelNode partial = tryPartialResultAggregate(aggregate, project, partitionFields); - if (partial != null) { - return partial; - } - } - } int queryBucketSize = osIndex.getQueryBucketSize(); boolean bucketNullable = !PPLHintUtils.ignoreNullBucket(aggregate); AggregateAnalyzer.AggregateBuilderHelper helper = @@ -457,111 +431,6 @@ private AbstractRelNode pushDownAggregate( return null; } - /** - * Resolve the aggregation's group keys to the underlying scan fields to partition indices on. A - * key may be a bare field ({@code ... by city}) or an expression over fields ({@code eval g = - * lower(city) | ... by g}); in the latter case we partition on every field the expression reads, - * since a kept index must map all of them aggregatably. The {@code project} (when present) sits - * directly on the scan, so its input refs index into this scan's row type. Returns {@code null} - * if any key is a pure constant with no field to key on. - */ - @Nullable - private List resolvePartitionFields(Aggregate aggregate, @Nullable Project project) { - List scanFields = getRowType().getFieldNames(); - List fields = new ArrayList<>(); - for (int group : aggregate.getGroupSet()) { - Set refs = new LinkedHashSet<>(); - if (project == null) { - refs.add(group); // group key indexes directly into the scan - } else { - project - .getProjects() - .get(group) - .accept( - new RexVisitorImpl(true) { - @Override - public Void visitInputRef(RexInputRef ref) { - refs.add(ref.getIndex()); - return null; - } - }); - } - if (refs.isEmpty()) { - return null; // constant group key -> nothing to partition on - } - for (int ref : refs) { - String name = scanFields.get(ref); - if (!fields.contains(name)) { - fields.add(name); - } - } - } - return fields; - } - - /** - * On a text/keyword mapping conflict, narrow the scan to the index subset where the group field - * is aggregatable, push the aggregation over just that subset, and record a warning naming the - * excluded indices. Only runs behind the opt-in setting and only when the response format can - * carry the warning ({@link CalcitePlanContext#isWarningsSupported}); returns {@code null} - * otherwise. {@code partitionFields} are the scan fields the group keys resolve to (see {@link - * #resolvePartitionFields}). Partitioning lives in {@link PartialResultAggregatePushdown}. - */ - private AbstractRelNode tryPartialResultAggregate( - Aggregate aggregate, @Nullable Project project, List partitionFields) { - // The per-request override wins when present; otherwise the cluster setting decides. Both the - // override and warnings-support are read from CalcitePlanContext (carried onto the plan), not - // Log4j ThreadContext, so they survive the security transport→worker handoff. - Boolean override = CalcitePlanContext.getPartialResultOverride(); - boolean partialResultEnabled = - override != null - ? override - : osIndex - .getSettings() - .getSettingValue(Settings.Key.PARTIAL_RESULT_ON_MAPPING_CONFLICT); - if (!partialResultEnabled) { - return null; - } - // A format with no warnings channel (CSV/RAW/VIZ) must not silently drop indices. - if (!CalcitePlanContext.isWarningsSupported()) { - return null; - } - try { - Map mappings = osIndex.getIndexMappings(); - PartialResultAggregatePushdown.Plan plan = - PartialResultAggregatePushdown.plan(partitionFields, mappings); - if (plan == null) { - return null; - } - - OpenSearchIndex narrowedIndex = - new OpenSearchIndex( - osIndex.getClient(), osIndex.getSettings(), String.join(",", plan.keptIndices())); - CalciteLogicalIndexScan narrowedScan = - new CalciteLogicalIndexScan( - getCluster(), - traitSet, - hints, - table, - narrowedIndex, - getRowType(), - pushDownContext.cloneWithOsIndex(narrowedIndex)); - // allowPartialFallback=false: the subset is already narrowed, so keep this one-shot. - AbstractRelNode pushed = narrowedScan.pushDownAggregate(aggregate, project, false); - if (pushed == null) { - return null; // narrowed subset still can't push down -> leave un-pushed - } - - CalcitePlanContext.addWarning(plan.warning()); - return pushed; - } catch (Exception e) { - if (LOG.isDebugEnabled()) { - LOG.debug("Cannot apply partial-result aggregate pushdown for {}", aggregate, e); - } - return null; - } - } - public AbstractRelNode pushDownLimit(LogicalSort sort, Integer limit, Integer offset) { try { if (pushDownContext.isAggregatePushed()) { diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/PartialResultAggregatePushdown.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/PartialResultAggregatePushdown.java deleted file mode 100644 index 460840f9cd2..00000000000 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/PartialResultAggregatePushdown.java +++ /dev/null @@ -1,157 +0,0 @@ -/* - * Copyright OpenSearch Contributors - * SPDX-License-Identifier: Apache-2.0 - */ - -package org.opensearch.sql.opensearch.storage.scan; - -import java.util.ArrayList; -import java.util.LinkedHashMap; -import java.util.List; -import java.util.Map; -import javax.annotation.Nullable; -import org.opensearch.sql.executor.Warning; -import org.opensearch.sql.opensearch.data.type.OpenSearchDataType; -import org.opensearch.sql.opensearch.data.type.OpenSearchDataType.MappingType; -import org.opensearch.sql.opensearch.mapping.IndexMapping; - -/** - * Selects which indices to keep so an aggregation on a mapping conflict can still push down over a - * clean subset. When a group field is mapped inconsistently across a wildcard pattern, some indices - * map it to a non-aggregatable type -- the text family ({@code text}, {@code text} with a {@code - * .keyword} sub-field, {@code match_only_text}), which the type merge collapses to bare {@code - * text} with no doc values -- while others map it to an aggregatable type ({@code keyword}, a - * numeric, {@code date}, {@code boolean}, {@code ip}). This drops the non-aggregatable indices, - * keeps the aggregatable ones, and attaches a warning naming what was excluded; the caller ({@link - * CalciteLogicalIndexScan}) re-runs pushdown over {@link Plan#keptIndices()}. - * - *

It applies only when the kept indices share one aggregatable type. If they would mix - * incompatible aggregatable types ({@code keyword} vs {@code integer}, two numeric types) the - * merged type is an arbitrary last-write-wins, so no subset can be kept safely; that is left to the - * normal path (a fundamental type conflict tracked separately as #5610). - */ -final class PartialResultAggregatePushdown { - - /** Max excluded index names to spell out in the warning; the rest are summarized as "N more". */ - static final int MAX_EXCLUDED_INDICES_IN_WARNING = 5; - - private PartialResultAggregatePushdown() {} - - /** - * The outcome of partitioning: which indices to aggregate over and the warning to attach. Absent - * (see {@link #plan}) when partial mode cannot or need not apply. - */ - record Plan(List keptIndices, List excludedIndices, Warning warning) {} - - /** - * Decide the partial-result plan for a group key over a set of per-index mappings. - * - * @param bucketNames the storage fields the group keys resolve to (dotted paths); an expression - * key like {@code lower(city)} resolves to the field(s) it reads, e.g. {@code city} - * @param mappings per-index field mappings, keyed by concrete index name (from {@code - * getIndexMappings}); the wildcard has already been resolved to concrete indices - * @return a plan naming the kept and excluded indices plus the warning, or {@code null} when - * partial mode does not apply: fewer than two indices, no aggregatable subset, the kept - * indices would mix incompatible aggregatable types, or nothing is excluded - */ - @Nullable - static Plan plan(List bucketNames, Map mappings) { - if (bucketNames.isEmpty() || mappings.size() < 2) { - return null; - } - - // Group indices by their aggregatable "compatibility signature": indices sharing a signature - // map every group field to the same aggregatable type, so they can be aggregated together - // without re-introducing a conflict. A null signature means the field is non-aggregatable in - // that index (text family or absent) -- always excludable. - Map> aggregatableGroups = new LinkedHashMap<>(); - List excludedIndices = new ArrayList<>(); - for (Map.Entry entry : mappings.entrySet()) { - // Flatten so a nested object field (mapping tree resource -> attributes -> applicationid) is - // keyed by its dotted path, matching the bucket field name Calcite resolved. - Map flatMapping = - OpenSearchDataType.traverseAndFlatten(entry.getValue().getFieldMappings()); - String signature = resolveBucketSignature(flatMapping, bucketNames); - if (signature == null) { - excludedIndices.add(entry.getKey()); - } else { - aggregatableGroups.computeIfAbsent(signature, k -> new ArrayList<>()).add(entry.getKey()); - } - } - - // Keep the aggregatable indices only when they share exactly one type. Zero aggregatable groups - // means partial mode can't help; more than one means incompatible aggregatable types whose - // merged type is arbitrary -> leave it to the normal path. - if (aggregatableGroups.size() != 1) { - return null; - } - if (excludedIndices.isEmpty()) { - return null; // homogeneous already -> pushdown would not have failed - } - - List keptIndices = aggregatableGroups.values().iterator().next(); - return new Plan( - keptIndices, excludedIndices, buildWarning(bucketNames, excludedIndices, mappings.size())); - } - - /** - * A per-index compatibility signature for the grouped field(s): one token per field joined by - * {@code |}. Two indices with equal signatures map every group field to the same aggregatable - * type and can be aggregated together. Returns {@code null} if any field is non-aggregatable - * here: absent, or a text-family type. A {@code text} field is non-aggregatable even with a - * {@code .keyword} sub-field, because the type merge collapses it to bare {@code text}. Each - * aggregatable field contributes a {@code t:TYPE} token (e.g. {@code t:keyword}, {@code - * t:integer}). - */ - @Nullable - static String resolveBucketSignature( - Map flatMapping, List bucketNames) { - List tokens = new ArrayList<>(); - for (String field : bucketNames) { - OpenSearchDataType type = flatMapping.get(field); - if (type == null) { - return null; // field absent here -> not aggregatable - } - MappingType mappingType = type.getMappingType(); - if (mappingType == MappingType.Text || mappingType == MappingType.MatchOnlyText) { - return null; // text family (incl. text-with-.keyword) collapses to bare text on merge - } - tokens.add("t:" + mappingType); // aggregatable type (keyword, numeric, date, boolean, ip) - } - return String.join("|", tokens); - } - - private static Warning buildWarning( - List bucketNames, List excludedIndices, int totalIndices) { - // Sort here (not in plan): ordering only matters for a stable, readable message. - List sortedExcluded = new ArrayList<>(excludedIndices); - sortedExcluded.sort(null); - String message = - String.format( - "Results exclude %d of %d indices due to a mapping conflict on %s.", - sortedExcluded.size(), totalIndices, bucketNames); - String detail = - String.format( - "%s is not aggregatable in every queried index (mapped as text or otherwise without doc" - + " values there), so these indices were excluded from the aggregation: %s. Map %s" - + " as an aggregatable type across all indices to include them.", - bucketNames, - formatIndexList(sortedExcluded, MAX_EXCLUDED_INDICES_IN_WARNING), - bucketNames); - return new Warning(Warning.TYPE_PARTIAL_RESULT, message, detail); - } - - /** - * Format a (sorted) index-name list for a warning: spell out up to {@code limit} names, then - * summarize any remainder as "and N more" so the message stays readable when the excluded set is - * large. - */ - static String formatIndexList(List indices, int limit) { - if (indices.size() <= limit) { - return indices.toString(); - } - return String.format( - "[%s, ... and %d more]", - String.join(", ", indices.subList(0, limit)), indices.size() - limit); - } -} diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/context/PushDownContext.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/context/PushDownContext.java index ba266a0f846..a622f948efb 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/context/PushDownContext.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/context/PushDownContext.java @@ -56,21 +56,6 @@ public PushDownContext clone() { return newContext; } - /** - * Clone this context but rebind it to a different {@link OpenSearchIndex}. Used by partial-result - * pushdown, which re-targets an aggregation at a narrowed index subset while preserving the - * operations already pushed onto the current scan (e.g. a WHERE filter). The pushed operations - * are index-agnostic request-builder actions, so replaying them onto the new index is safe. - */ - public PushDownContext cloneWithOsIndex(OpenSearchIndex newOsIndex) { - PushDownContext newContext = new PushDownContext(newOsIndex); - for (PushDownOperation operation : this) { - newContext.add(operation); - } - newContext.aggSpec = aggSpec; - return newContext; - } - /** * Create a new {@link PushDownContext} without the collation action. * diff --git a/opensearch/src/test/java/org/opensearch/sql/opensearch/request/system/OpenSearchDescribeIndexRequestTest.java b/opensearch/src/test/java/org/opensearch/sql/opensearch/request/system/OpenSearchDescribeIndexRequestTest.java index ef2af9cd9a7..6e156e6d333 100644 --- a/opensearch/src/test/java/org/opensearch/sql/opensearch/request/system/OpenSearchDescribeIndexRequestTest.java +++ b/opensearch/src/test/java/org/opensearch/sql/opensearch/request/system/OpenSearchDescribeIndexRequestTest.java @@ -36,14 +36,12 @@ class OpenSearchDescribeIndexRequestTest { @Mock private IndexMapping mapping2; /** - * Merging must not mutate the per-index mappings it reads. {@code MergeRuleHelper} rewrites the - * accumulated type's nested {@code properties} in place, so without copying, the first index's - * nested field would be merged into the second's -- and callers of {@link - * OpenSearchDescribeIndexRequest#getLastIndexMappings()} (partial-result partitioning) would see - * a text/keyword conflict as no conflict at all. + * Merging must not mutate the mappings it reads: {@code MergeRuleHelper} rewrites the accumulated + * type's nested {@code properties} in place, so without a copy the first index's nested field + * would be merged into the second's and a fetched mapping would no longer describe its index. */ @Test - void getFieldTypesLeavesRetainedPerIndexMappingsIntact() { + void getFieldTypesLeavesFetchedMappingsIntact() { Map keywordSide = Map.of("attrs", nestedObject("env", OpenSearchDataType.MappingType.Keyword)); Map textSide = @@ -56,13 +54,10 @@ void getFieldTypesLeavesRetainedPerIndexMappingsIntact() { OpenSearchDescribeIndexRequest request = new OpenSearchDescribeIndexRequest(client, "idx-*"); request.getFieldTypes(); - // Each index must still report the type it actually declared. + // Each fetched mapping must still declare the type it came with. assertEquals( - OpenSearchDataType.MappingType.Keyword, - nestedFieldType(request.getLastIndexMappings().get("idx-keyword"), "attrs", "env")); - assertEquals( - OpenSearchDataType.MappingType.Text, - nestedFieldType(request.getLastIndexMappings().get("idx-text"), "attrs", "env")); + OpenSearchDataType.MappingType.Keyword, nestedFieldType(keywordSide, "attrs", "env")); + assertEquals(OpenSearchDataType.MappingType.Text, nestedFieldType(textSide, "attrs", "env")); } private static OpenSearchDataType nestedObject( @@ -104,8 +99,8 @@ private static OpenSearchDataType nestedObjectTo(String leafField, Map mappings, String parent, String child) { + return mappings.get(parent).getProperties().get(child).getMappingType(); } @Test diff --git a/opensearch/src/test/java/org/opensearch/sql/opensearch/storage/scan/PartialResultAggregatePushdownTest.java b/opensearch/src/test/java/org/opensearch/sql/opensearch/storage/scan/PartialResultAggregatePushdownTest.java deleted file mode 100644 index 6a269269363..00000000000 --- a/opensearch/src/test/java/org/opensearch/sql/opensearch/storage/scan/PartialResultAggregatePushdownTest.java +++ /dev/null @@ -1,254 +0,0 @@ -/* - * Copyright OpenSearch Contributors - * SPDX-License-Identifier: Apache-2.0 - */ - -package org.opensearch.sql.opensearch.storage.scan; - -import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.assertNull; -import static org.junit.jupiter.api.Assertions.assertTrue; - -import java.util.LinkedHashMap; -import java.util.List; -import java.util.Map; -import java.util.stream.Collectors; -import java.util.stream.IntStream; -import org.junit.jupiter.api.Test; -import org.opensearch.sql.executor.Warning; -import org.opensearch.sql.opensearch.data.type.OpenSearchDataType; -import org.opensearch.sql.opensearch.data.type.OpenSearchDataType.MappingType; -import org.opensearch.sql.opensearch.data.type.OpenSearchTextType; -import org.opensearch.sql.opensearch.mapping.IndexMapping; -import org.opensearch.sql.opensearch.storage.scan.PartialResultAggregatePushdown.Plan; - -class PartialResultAggregatePushdownTest { - - private static final OpenSearchDataType KEYWORD_TYPE = OpenSearchDataType.of(MappingType.Keyword); - private static final OpenSearchDataType BARE_TEXT_TYPE = OpenSearchTextType.of(); - private static final OpenSearchDataType INT_TYPE = OpenSearchDataType.of(MappingType.Integer); - private static final OpenSearchDataType TEXT_WITH_KEYWORD_TYPE = - OpenSearchTextType.of(Map.of("keyword", OpenSearchDataType.of(MappingType.Keyword))); - - // ---- resolveBucketSignature ---- - - @Test - void resolveKeywordField() { - assertEquals( - "t:keyword", resolveOne(Map.of("f", KEYWORD_TYPE)), "keyword resolves to its type token"); - } - - @Test - void resolveBareTextFieldIsNonAggregatable() { - assertNull(resolveOne(Map.of("f", BARE_TEXT_TYPE)), "bare text is not aggregatable"); - } - - @Test - void resolveTextWithKeywordFieldIsNonAggregatable() { - // A text field with a .keyword sub-field still merges to bare text across a conflict, so for - // partial-result purposes it is non-aggregatable, same as bare text. - assertNull( - resolveOne(Map.of("f", TEXT_WITH_KEYWORD_TYPE)), - "text with a .keyword sub-field is non-aggregatable (collapses to text on merge)"); - } - - @Test - void resolveNonTextTypeIsAggregatable() { - assertEquals( - "t:integer", - resolveOne(Map.of("f", INT_TYPE)), - "a numeric field is aggregatable, tagged by its concrete type"); - } - - @Test - void resolveAbsentField() { - assertNull( - PartialResultAggregatePushdown.resolveBucketSignature(Map.of(), List.of("f")), - "a field absent from the index cannot be aggregated cleanly"); - } - - @Test - void resolveMultiFieldSignatureIsPerFieldTokens() { - Map mapping = Map.of("a", KEYWORD_TYPE, "b", INT_TYPE); - assertEquals( - "t:keyword|t:integer", - PartialResultAggregatePushdown.resolveBucketSignature(mapping, List.of("a", "b"))); - // A text field anywhere in the key makes the whole index non-aggregatable. - Map withText = - Map.of("a", KEYWORD_TYPE, "b", INT_TYPE, "c", BARE_TEXT_TYPE); - assertNull( - PartialResultAggregatePushdown.resolveBucketSignature(withText, List.of("a", "b", "c"))); - } - - // ---- plan: not-applicable cases return null ---- - - @Test - void planNullForSingleIndex() { - assertNull(PartialResultAggregatePushdown.plan(List.of("f"), Map.of("idx", keywordIndex()))); - } - - @Test - void planNullForEmptyBucketNames() { - assertNull( - PartialResultAggregatePushdown.plan( - List.of(), Map.of("kw", keywordIndex(), "txt", bareTextIndex()))); - } - - @Test - void planNullWhenNoConflict() { - // Two keyword indices -> one signature, nothing excluded -> pushdown would have worked. - assertNull( - PartialResultAggregatePushdown.plan( - List.of("f"), Map.of("kw1", keywordIndex(), "kw2", keywordIndex()))); - } - - @Test - void planNullWhenNoAggregatableSubset() { - // Every index is text (bare or text-with-.keyword) -> nothing aggregatable to keep. - assertNull( - PartialResultAggregatePushdown.plan( - List.of("f"), Map.of("txt", bareTextIndex(), "tk", textWithKeywordIndex()))); - } - - @Test - void planNullWhenKeptWouldMixIncompatibleTypes() { - // keyword vs int: two aggregatable signatures, arbitrary merged type -> bail (#5610). - assertNull( - PartialResultAggregatePushdown.plan( - List.of("f"), ordered("kw", keywordIndex(), "num", intIndex()))); - // int vs long: two numeric signatures -> bail. - assertNull( - PartialResultAggregatePushdown.plan( - List.of("f"), ordered("i", intIndex(), "l", longIndex()))); - // keyword + int + text: still mixes keyword and int, even though the text index alone is - // excludable -> bail. - assertNull( - PartialResultAggregatePushdown.plan( - List.of("f"), - ordered("kw", keywordIndex(), "num", intIndex(), "txt", bareTextIndex()))); - } - - // ---- plan: partitioning ---- - - @Test - void planKeepsKeywordExcludesBareText() { - Plan plan = - PartialResultAggregatePushdown.plan( - List.of("f"), ordered("kw", keywordIndex(), "txt", bareTextIndex())); - assertEquals(List.of("kw"), plan.keptIndices()); - assertEquals(List.of("txt"), plan.excludedIndices()); - assertEquals(Warning.TYPE_PARTIAL_RESULT, plan.warning().getType()); - } - - @Test - void planKeepsKeywordExcludesTextWithKeyword() { - // text-with-.keyword is non-aggregatable for partial purposes, so keyword is kept and the - // text-with-.keyword index is excluded. - Plan plan = - PartialResultAggregatePushdown.plan( - List.of("f"), ordered("kw", keywordIndex(), "tk", textWithKeywordIndex())); - assertEquals(List.of("kw"), plan.keptIndices()); - assertEquals(List.of("tk"), plan.excludedIndices()); - } - - @Test - void planKeepsNumericExcludesText() { - // int vs bare text: int is aggregatable, text is not -> keep the int index, exclude text. - Plan plan = - PartialResultAggregatePushdown.plan( - List.of("f"), ordered("num", intIndex(), "txt", bareTextIndex())); - assertEquals(List.of("num"), plan.keptIndices()); - assertEquals(List.of("txt"), plan.excludedIndices()); - } - - @Test - void planKeepsAllKeywordIndicesExcludesText() { - // Two keyword indices + one text index -> keep both keyword, exclude the text one. - Plan plan = - PartialResultAggregatePushdown.plan( - List.of("f"), - ordered("kw1", keywordIndex(), "kw2", keywordIndex(), "txt", bareTextIndex())); - assertEquals(List.of("kw1", "kw2"), plan.keptIndices()); - assertEquals(List.of("txt"), plan.excludedIndices()); - } - - @Test - void warningListsExcludedIndicesSorted() { - Map mappings = - ordered( - "kw", keywordIndex(), - "txt-c", bareTextIndex(), - "txt-a", bareTextIndex(), - "txt-b", bareTextIndex()); - Plan plan = PartialResultAggregatePushdown.plan(List.of("f"), mappings); - // The plan keeps insertion order; only the warning message sorts them for readability. - assertTrue(plan.warning().getDetail().contains("[txt-a, txt-b, txt-c]")); - } - - // ---- formatIndexList ---- - - @Test - void formatShortListInFull() { - assertEquals( - "[a, b, c]", PartialResultAggregatePushdown.formatIndexList(List.of("a", "b", "c"), 5)); - } - - @Test - void formatLongListTruncated() { - List many = - IntStream.rangeClosed(1, 8).mapToObj(i -> "idx" + i).collect(Collectors.toList()); - String formatted = PartialResultAggregatePushdown.formatIndexList(many, 5); - assertTrue(formatted.contains("idx1"), "spells out the first few"); - assertTrue(formatted.contains("idx5"), "spells out up to the cap"); - assertTrue(formatted.contains("and 3 more"), "summarizes the remainder"); - assertTrue(!formatted.contains("idx6"), "does not list beyond the cap"); - } - - // ---- warning content ---- - - @Test - void warningNamesFieldExcludedIndicesAndCount() { - Plan plan = - PartialResultAggregatePushdown.plan( - List.of("f"), ordered("kw", keywordIndex(), "txt", bareTextIndex())); - Warning w = plan.warning(); - assertTrue(w.getMessage().contains("1 of 2 indices"), "message reports the excluded count"); - assertTrue(w.getMessage().contains("f"), "message names the field"); - assertTrue(w.getDetail().contains("txt"), "detail names the excluded index"); - } - - // ---- helpers ---- - - private static String resolveOne(Map mapping) { - return PartialResultAggregatePushdown.resolveBucketSignature(mapping, List.of("f")); - } - - private static IndexMapping keywordIndex() { - return new IndexMapping(Map.of("f", KEYWORD_TYPE)); - } - - private static IndexMapping bareTextIndex() { - return new IndexMapping(Map.of("f", BARE_TEXT_TYPE)); - } - - private static IndexMapping textWithKeywordIndex() { - return new IndexMapping(Map.of("f", TEXT_WITH_KEYWORD_TYPE)); - } - - private static IndexMapping intIndex() { - return new IndexMapping(Map.of("f", INT_TYPE)); - } - - private static IndexMapping longIndex() { - return new IndexMapping(Map.of("f", OpenSearchDataType.of(MappingType.Long))); - } - - /** Build an insertion-ordered map so kept/excluded assertions are deterministic. */ - private static Map ordered(Object... keyThenValue) { - Map map = new LinkedHashMap<>(); - for (int i = 0; i < keyThenValue.length; i += 2) { - map.put((String) keyThenValue[i], (IndexMapping) keyThenValue[i + 1]); - } - return map; - } -} diff --git a/plugin/src/main/java/org/opensearch/sql/plugin/request/PPLQueryRequestFactory.java b/plugin/src/main/java/org/opensearch/sql/plugin/request/PPLQueryRequestFactory.java index 1f42b071621..8d13e42b8e4 100644 --- a/plugin/src/main/java/org/opensearch/sql/plugin/request/PPLQueryRequestFactory.java +++ b/plugin/src/main/java/org/opensearch/sql/plugin/request/PPLQueryRequestFactory.java @@ -33,7 +33,6 @@ public class PPLQueryRequestFactory { private static final String QUERY_PARAMS_ANALYZE = "analyze"; private static final String QUERY_PARAMS_FETCH_SIZE = "fetch_size"; private static final String QUERY_PARAMS_INCLUDE_METADATA = "include_metadata"; - private static final String QUERY_PARAMS_PARTIAL_RESULT = "partial_result"; /** * Build {@link PPLQueryRequest} from {@link RestRequest}. @@ -132,11 +131,6 @@ private static PPLQueryRequest parsePPLRequestFromPayload(RestRequest restReques if (queryId != null) { pplRequest.queryId(queryId); } - // Set the override only when present, so a request that omits it defers to the cluster - // setting. - if (jsonContent.has(QUERY_PARAMS_PARTIAL_RESULT)) { - pplRequest.partialResult(jsonContent.optBoolean(QUERY_PARAMS_PARTIAL_RESULT)); - } return pplRequest; } catch (JSONException e) { throw new IllegalArgumentException("Failed to parse request payload", e); diff --git a/plugin/src/main/java/org/opensearch/sql/plugin/transport/TransportPPLQueryAction.java b/plugin/src/main/java/org/opensearch/sql/plugin/transport/TransportPPLQueryAction.java index 5a16c60b99f..6754995ea06 100644 --- a/plugin/src/main/java/org/opensearch/sql/plugin/transport/TransportPPLQueryAction.java +++ b/plugin/src/main/java/org/opensearch/sql/plugin/transport/TransportPPLQueryAction.java @@ -194,13 +194,6 @@ protected void doExecute( // in order to use PPL service, we need to convert TransportPPLQueryRequest to PPLQueryRequest PPLQueryRequest transformedRequest = transportRequest.toPPLQueryRequest(); QueryContext.setProfile(transformedRequest.profile()); - // Only the JSON shape carries warnings; gate partial results on it so CSV/RAW/VIZ never drop - // data silently. Carried on the request (not Log4j ThreadContext) so it survives the - // transport→worker handoff, which the security plugin's interceptor does not preserve. - transformedRequest.warningsSupported(warningsSupported(transformedRequest)); - // The per-request partial-result override (e.g. a Dashboards toggle) rides on the request → - // plan → worker thread (see PPLService/QueryPlan), not Log4j ThreadContext, for the same - // handoff-survival reason as warningsSupported. null defers to the cluster setting. // Start root span with OTel DB semantic convention attributes Span rootSpan = @@ -406,22 +399,6 @@ private Format format(PPLQueryRequest pplRequest) { } } - /** - * Whether the requested response format carries a warnings channel. Only the JSON shape (built by - * {@code SimpleJsonResponseFormatter} -- the fallback for anything that is not CSV/RAW/VIZ) emits - * warnings; the others have no slot for them. Mirrors the format branching in {@link - * #createListener}. Explain requests are excluded up front: their {@code format} is an - * explain-only value (e.g. {@code json}/{@code yaml}) that {@link #format} cannot resolve, and an - * explain response never carries query warnings. - */ - private boolean warningsSupported(PPLQueryRequest pplRequest) { - if (pplRequest.isExplainRequest()) { - return false; - } - Format format = format(pplRequest); - return !(format.equals(Format.CSV) || format.equals(Format.RAW) || format.equals(Format.VIZ)); - } - private ActionListener wrapWithProfilingClear( ActionListener delegate) { return new ActionListener<>() { @@ -447,9 +424,7 @@ public void onFailure(Exception e) { /** * Clear the per-request state carried in {@link QueryContext}'s thread-locals. Transport threads - * are pooled, so anything left behind is inherited by the next query to run on this thread. (The - * partial-result override no longer lives here -- it rides on the plan to the worker thread and - * is reset per query in {@code CalcitePlanContext}.) + * are pooled, so anything left behind is inherited by the next query to run on this thread. */ private static void clearRequestScopedState() { QueryProfiling.clear(); diff --git a/plugin/src/main/java/org/opensearch/sql/plugin/transport/TransportPPLQueryRequest.java b/plugin/src/main/java/org/opensearch/sql/plugin/transport/TransportPPLQueryRequest.java index 13c0f579ab3..68a96a4f924 100644 --- a/plugin/src/main/java/org/opensearch/sql/plugin/transport/TransportPPLQueryRequest.java +++ b/plugin/src/main/java/org/opensearch/sql/plugin/transport/TransportPPLQueryRequest.java @@ -63,15 +63,6 @@ public class TransportPPLQueryRequest extends ActionRequest { @Accessors(fluent = true) private String queryId = null; - /** - * Per-request override for partial-result mode; null means defer to the cluster setting. See - * {@link PPLQueryRequest#partialResult()}. - */ - @Setter - @Getter - @Accessors(fluent = true) - private Boolean partialResult = null; - /** Constructor of TransportPPLQueryRequest from PPLQueryRequest. */ public TransportPPLQueryRequest(PPLQueryRequest pplQueryRequest) { pplQuery = pplQueryRequest.getRequest(); @@ -84,7 +75,6 @@ public TransportPPLQueryRequest(PPLQueryRequest pplQueryRequest) { analyze = pplQueryRequest.analyze(); explainMode = pplQueryRequest.mode().getModeName(); queryId = pplQueryRequest.queryId(); - partialResult = pplQueryRequest.partialResult(); } /** Constructor of TransportPPLQueryRequest from StreamInput. */ @@ -101,7 +91,6 @@ public TransportPPLQueryRequest(StreamInput in) throws IOException { profile = in.readBoolean(); analyze = in.readBoolean(); queryId = in.readOptionalString(); - partialResult = in.readOptionalBoolean(); } /** Re-create the object from the actionRequest. */ @@ -136,7 +125,6 @@ public void writeTo(StreamOutput out) throws IOException { out.writeBoolean(profile); out.writeBoolean(analyze); out.writeOptionalString(queryId); - out.writeOptionalBoolean(partialResult); } public String getRequest() { @@ -196,7 +184,6 @@ public PPLQueryRequest toPPLQueryRequest() { pplQueryRequest.sanitize(sanitize); pplQueryRequest.style(style); pplQueryRequest.queryId(queryId); - pplQueryRequest.partialResult(partialResult); return pplQueryRequest; } } diff --git a/ppl/src/main/java/org/opensearch/sql/ppl/PPLService.java b/ppl/src/main/java/org/opensearch/sql/ppl/PPLService.java index e2572b5f0da..8d82b550d44 100644 --- a/ppl/src/main/java/org/opensearch/sql/ppl/PPLService.java +++ b/ppl/src/main/java/org/opensearch/sql/ppl/PPLService.java @@ -211,8 +211,6 @@ private AbstractPlan plan( anonymizedQuerySink.accept(anonymized); AbstractPlan plan = queryExecutionFactory.create(statement, queryListener, explainListener); - plan.setWarningsSupported(request.warningsSupported()); - plan.setPartialResultOverride(request.partialResult()); return plan; } } diff --git a/ppl/src/main/java/org/opensearch/sql/ppl/domain/PPLQueryRequest.java b/ppl/src/main/java/org/opensearch/sql/ppl/domain/PPLQueryRequest.java index 67fff2db5e5..39a3ce2258e 100644 --- a/ppl/src/main/java/org/opensearch/sql/ppl/domain/PPLQueryRequest.java +++ b/ppl/src/main/java/org/opensearch/sql/ppl/domain/PPLQueryRequest.java @@ -73,25 +73,6 @@ public class PPLQueryRequest { @Accessors(fluent = true) private String queryId = null; - /** - * Per-request override for partial-result mode. {@code null} means the request expressed no - * preference and the cluster setting decides; non-null forces partial mode on/off for this query. - */ - @Setter - @Getter - @Accessors(fluent = true) - private Boolean partialResult = null; - - /** - * Whether the requested response format can carry a warnings channel. Derived from the format on - * the transport thread and threaded to the worker via the plan, so the partial-result gate does - * not rely on Log4j ThreadContext (dropped across the security plugin's thread handoff). - */ - @Setter - @Getter - @Accessors(fluent = true) - private boolean warningsSupported = false; - public PPLQueryRequest(String pplQuery, JSONObject jsonContent, String path) { this(pplQuery, jsonContent, path, ""); }