diff --git a/.github/workflows/codeql-analysis.yml b/.github/workflows/codeql-analysis.yml
index 2186c166e19a6..e9552e0537f9b 100644
--- a/.github/workflows/codeql-analysis.yml
+++ b/.github/workflows/codeql-analysis.yml
@@ -51,6 +51,9 @@ jobs:
# Prefix the list here with "+" to use these queries and those in the config file.
# queries: ./path/to/local/query, your-org/your-repo/queries@main
+ - name: Set up protoc
+ uses: arduino/setup-protoc@v3
+
# Autobuild attempts to build any compiled languages (C/C++, C#, or Java).
# If this step fails, then you should remove it and run the build manually (see below)
- name: Autobuild
diff --git a/sandbox/libs/analytics-framework/src/main/java/org/opensearch/analytics/spi/AnalyticsSearchBackendPlugin.java b/sandbox/libs/analytics-framework/src/main/java/org/opensearch/analytics/spi/AnalyticsSearchBackendPlugin.java
index c823763e2040d..373fcde77b75f 100644
--- a/sandbox/libs/analytics-framework/src/main/java/org/opensearch/analytics/spi/AnalyticsSearchBackendPlugin.java
+++ b/sandbox/libs/analytics-framework/src/main/java/org/opensearch/analytics/spi/AnalyticsSearchBackendPlugin.java
@@ -8,6 +8,9 @@
package org.opensearch.analytics.spi;
+import org.opensearch.analytics.backend.EngineResultStream;
+import org.opensearch.analytics.backend.ExecutionContext;
+import org.opensearch.analytics.backend.SearchExecEngine;
import org.opensearch.index.engine.dataformat.DataFormat;
import java.util.Collections;
@@ -17,9 +20,26 @@
/**
* SPI extension point for back-end query engines for query planning and execution capabilities
* as needed by the {@link org.opensearch.analytics.exec.QueryPlanExecutor}
+ *
+ *
TODO: separate capability declaration (planner, coordinator) from execution engine factory
+ * (data node) into two interfaces. AnalyticsSearchBackendPlugin should only declare capabilities.
+ * SearchExecEngineProvider should be discovered separately by the executor. Remove the extends
+ * relationship and the default createSearchExecEngine() below once that separation is done.
*/
public interface AnalyticsSearchBackendPlugin extends SearchExecEngineProvider {
+ /** Unique engine name (e.g., "lucene", "datafusion"). */
+ String name();
+
+ /**
+ * {@inheritDoc}
+ * Temporary default — remove once SearchExecEngineProvider is separated from this interface.
+ */
+ @Override
+ default SearchExecEngine createSearchExecEngine(ExecutionContext ctx) {
+ throw new UnsupportedOperationException("createSearchExecEngine not implemented for " + name());
+ }
+
/** Returns the data formats supported by this backend. */
List getSupportedFormats();
@@ -71,5 +91,4 @@ default Set supportedShuffleCapabilities() {
default byte[] convertFragment(Object fragment) {
throw new UnsupportedOperationException("convertFragment not yet implemented for " + name());
}
-
}
diff --git a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/exec/DefaultPlanExecutor.java b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/exec/DefaultPlanExecutor.java
index 1875b2fbb49a3..e1d42f9f5a645 100644
--- a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/exec/DefaultPlanExecutor.java
+++ b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/exec/DefaultPlanExecutor.java
@@ -69,7 +69,6 @@ public Iterable execute(RelNode logicalFragment, Object context) {
new PlannerContext(
new CapabilityRegistry(new ArrayList<>(backEnds.values())),
clusterService.state()));
-
String tableName = extractTableName(logicalFragment);
AnalyticsSearchBackendPlugin provider = selectBackEnd();
if (provider == null) {
diff --git a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/FieldStorageResolver.java b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/FieldStorageResolver.java
index bfef4645ce79b..4858ce1d788c8 100644
--- a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/FieldStorageResolver.java
+++ b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/FieldStorageResolver.java
@@ -39,6 +39,21 @@ public class FieldStorageResolver {
/**
* Production: resolves per-field storage from IndexMetadata and backend capabilities.
+ *
+ * LIMITATION: This infers storage from what each DataFormat DECLARES it can support
+ * (via FieldTypeCapabilities), not from what the index ACTUALLY stores. A format declaring
+ * COLUMNAR_STORAGE for "integer" doesn't mean this index has integer doc values in that format.
+ *
+ *
The indexing side has no per-field format metadata at the coordinator level today:
+ * - DataFormat.supportedFields() = capability, not actual storage
+ * - Segment.dfGroupedSearchableFiles = per-segment format info, data node only
+ * - DataFormatRegistry has TODOs to filter by index settings/mapper service
+ * - primary_data_format index setting is the only coordinator-level hint (index-level, not field-level)
+ *
+ *
TODO: Replace inference with actual per-field format metadata once the indexing team adds
+ * it to MappingMetadata or IndexMetadata. Until then, this over-estimates viable backends —
+ * the shard-level cost function must handle the mismatch. Consider using primary_data_format
+ * as a narrowing hint in the interim.
*/
@SuppressWarnings("unchecked")
public FieldStorageResolver(IndexMetadata indexMetadata, List backends) {
diff --git a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/PlannerImpl.java b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/PlannerImpl.java
index f1ce9de977e7c..c2c7c902603a6 100644
--- a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/PlannerImpl.java
+++ b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/PlannerImpl.java
@@ -14,6 +14,8 @@
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.opensearch.analytics.planner.rel.OpenSearchDistributionTraitDef;
+import org.opensearch.analytics.planner.dag.DAGBuilder;
+import org.opensearch.analytics.planner.dag.QueryDAG;
import org.opensearch.analytics.planner.rules.OpenSearchAggregateRule;
import org.opensearch.analytics.planner.rules.OpenSearchAggregateSplitRule;
import org.opensearch.analytics.planner.rules.OpenSearchFilterRule;
@@ -21,22 +23,26 @@
import org.opensearch.analytics.planner.rules.OpenSearchSortRule;
import org.opensearch.analytics.planner.rules.OpenSearchTableScanRule;
+import org.opensearch.analytics.planner.dag.DAGBuilder;
+import org.opensearch.analytics.planner.dag.QueryDAG;
+
import java.util.List;
/**
* Central planner for the Analytics Plugin.
*
- * Two phases:
+ *
Three phases:
*
* HepPlanner (RBO): converts LogicalXxx → OpenSearchXxx with backend
* assignment, predicate annotation, and distribution traits.
* VolcanoPlanner (CBO): requests SINGLETON at root (coordinator must
* gather all results). Split rule fires on aggregates, Volcano inserts
* exchanges via trait enforcement where distribution mismatches.
+ * DAG construction: cuts at exchange boundaries, builds stage tree.
*
*
* TODO: eliminate copyToCluster — have frontends create RelNodes with Volcano cluster.
- *
TODO: DAG construction (cut at exchange boundaries)
+ *
TODO: Per-stage plan forking (multiple plan generation)
*
TODO: Fragment conversion (backend.convertFragment)
*
TODO: Join strategy selection, sort removal via CBO
*
@@ -93,6 +99,11 @@ public static RelNode createPlan(RelNode rawRelNode, PlannerContext context) {
RelNode result = volcanoPlanner.findBestExp();
LOGGER.info("After CBO:\n{}", RelOptUtil.toString(result));
+
+ // Phase 3: DAG construction — cut at exchange boundaries
+ QueryDAG dag = DAGBuilder.build(result);
+ LOGGER.info("QueryDAG:\n{}", dag);
+
return result;
}
}
diff --git a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/dag/DAGBuilder.java b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/dag/DAGBuilder.java
new file mode 100644
index 0000000000000..861c88a4c7109
--- /dev/null
+++ b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/dag/DAGBuilder.java
@@ -0,0 +1,159 @@
+/*
+ * SPDX-License-Identifier: Apache-2.0
+ *
+ * The OpenSearch Contributors require contributions made to
+ * this file be licensed under the Apache-2.0 license or a
+ * compatible open source license.
+ */
+
+package org.opensearch.analytics.planner.dag;
+
+import org.apache.calcite.rel.RelDistribution;
+import org.apache.calcite.rel.RelNode;
+import org.apache.calcite.rel.type.RelDataType;
+import org.opensearch.analytics.planner.rel.OpenSearchExchangeReducer;
+import org.opensearch.analytics.planner.rel.OpenSearchExchangeWriter;
+import org.opensearch.analytics.planner.rel.OpenSearchShuffleReader;
+import org.opensearch.analytics.planner.rel.OpenSearchStageInputScan;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.UUID;
+
+/**
+ * Builds a {@link QueryDAG} from the CBO output by cutting at exchange boundaries.
+ *
+ *
SINGLETON: {@link OpenSearchExchangeReducer} is the boundary. Parent
+ * fragment is everything above the reducer (may be null for a pure gather).
+ * Child fragment is the reducer's input subtree.
+ *
+ *
HASH/RANGE: {@link OpenSearchShuffleReader} → {@link OpenSearchExchangeWriter}.
+ * ShuffleReader stays in parent as leaf (input severed). Writer + subtree
+ * below becomes the child fragment.
+ *
+ *
Stage IDs assigned bottom-up (leaf stages get lower IDs).
+ *
+ * @opensearch.internal
+ */
+public class DAGBuilder {
+
+ private DAGBuilder() {}
+
+ public static QueryDAG build(RelNode cboOutput) {
+ int[] counter = { 0 };
+ List childStages = new ArrayList<>();
+
+ RelNode fragment;
+ if (cboOutput instanceof OpenSearchExchangeReducer reducer) {
+ // Root is an ExchangeReducer (e.g., shuffle case where coordinator
+ // is a pure gather). ExchangeReducer becomes the root stage fragment
+ // with its input severed. Child stage is the subtree below.
+ fragment = cutSingleton(reducer, counter, childStages);
+ } else {
+ fragment = sever(cboOutput, counter, childStages);
+ }
+
+ Stage rootStage = new Stage(counter[0]++, fragment, childStages, null);
+ return new QueryDAG(UUID.randomUUID().toString(), rootStage);
+ }
+
+ /**
+ * Walks top-down collecting operators into the current stage. When an
+ * exchange boundary is hit, severs the link: subtree below becomes a
+ * child stage, current stage continues with the boundary removed or
+ * replaced by a detached leaf.
+ *
+ * Exchange boundaries are detected at the input level (not the node
+ * level) so the parent node never receives a null input. For SINGLETON
+ * cuts the input is simply dropped; for shuffle cuts the ShuffleReader
+ * replaces the input as a detached leaf.
+ *
+ * @param node current node being visited
+ * @param counter stage ID counter (bottom-up assignment)
+ * @param childStages accumulator for child stages found below exchanges
+ * @return the rewritten fragment root for the current stage
+ */
+ private static RelNode sever(RelNode node, int[] counter, List childStages) {
+ // Check each input for exchange boundaries before recursing
+ List newInputs = new ArrayList<>();
+ for (RelNode input : node.getInputs()) {
+ if (input instanceof OpenSearchExchangeReducer reducer) {
+ // SINGLETON cut: ExchangeReducer stays in parent as leaf (input severed).
+ // Child stage is the reducer's input subtree. Analytics Core streams
+ // results from data nodes — no exchange operator needed in child fragment.
+ newInputs.add(cutSingleton(reducer, counter, childStages));
+ } else if (input instanceof OpenSearchShuffleReader reader) {
+ // Shuffle cut: ShuffleReader stays in parent as leaf (input severed).
+ // Child stage is ExchangeWriter + subtree below.
+ newInputs.add(cutShuffle(reader, counter, childStages));
+ } else {
+ newInputs.add(sever(input, counter, childStages));
+ }
+ }
+
+ // Leaf node (e.g., TableScan) — no inputs to process
+ if (node.getInputs().isEmpty()) {
+ return node;
+ }
+
+ // Rebuild only if inputs changed
+ boolean changed = false;
+ for (int idx = 0; idx < newInputs.size(); idx++) {
+ if (newInputs.get(idx) != node.getInputs().get(idx)) {
+ changed = true;
+ break;
+ }
+ }
+ return changed ? node.copy(node.getTraitSet(), newInputs) : node;
+ }
+
+ private static RelNode cutSingleton(OpenSearchExchangeReducer reducer,
+ int[] counter, List parentChildStages) {
+ List grandchildren = new ArrayList<>();
+ RelNode childFragment = sever(reducer.getInput(), counter, grandchildren);
+
+ int childStageId = counter[0]++;
+ parentChildStages.add(new Stage(
+ childStageId, childFragment, grandchildren,
+ new ExchangeInfo(RelDistribution.Type.SINGLETON, null, List.of())
+ ));
+
+ // Replace child subtree with StageInputScan, keep ExchangeReducer in parent
+ RelDataType childRowType = reducer.getInput().getRowType();
+ OpenSearchStageInputScan stageInput = new OpenSearchStageInputScan(
+ reducer.getCluster(), reducer.getTraitSet(), childStageId, childRowType
+ );
+ return new OpenSearchExchangeReducer(
+ reducer.getCluster(), reducer.getTraitSet(), stageInput, reducer.getViableBackends()
+ );
+ }
+
+ private static RelNode cutShuffle(OpenSearchShuffleReader reader,
+ int[] counter, List parentChildStages) {
+ if (!(reader.getInput() instanceof OpenSearchExchangeWriter writer)) {
+ throw new IllegalStateException(
+ "ShuffleReader input must be ExchangeWriter, got: "
+ + reader.getInput().getClass().getSimpleName());
+ }
+
+ List grandchildren = new ArrayList<>();
+ RelNode belowWriter = sever(writer.getInput(), counter, grandchildren);
+ RelNode childFragment = writer.copy(writer.getTraitSet(), List.of(belowWriter));
+
+ int childStageId = counter[0]++;
+ parentChildStages.add(new Stage(
+ childStageId, childFragment, grandchildren,
+ new ExchangeInfo(RelDistribution.Type.HASH_DISTRIBUTED, writer.getShuffleImpl(), writer.getKeys())
+ ));
+
+ // Replace child subtree with StageInputScan, keep ShuffleReader in parent
+ RelDataType childRowType = writer.getInput().getRowType();
+ OpenSearchStageInputScan stageInput = new OpenSearchStageInputScan(
+ reader.getCluster(), reader.getTraitSet(), childStageId, childRowType
+ );
+ return new OpenSearchShuffleReader(
+ reader.getCluster(), reader.getTraitSet(), stageInput,
+ reader.getViableBackends(), reader.getShuffleImpl()
+ );
+ }
+}
diff --git a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/dag/ExchangeInfo.java b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/dag/ExchangeInfo.java
new file mode 100644
index 0000000000000..d71ad6a017123
--- /dev/null
+++ b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/dag/ExchangeInfo.java
@@ -0,0 +1,37 @@
+/*
+ * SPDX-License-Identifier: Apache-2.0
+ *
+ * The OpenSearch Contributors require contributions made to
+ * this file be licensed under the Apache-2.0 license or a
+ * compatible open source license.
+ */
+
+package org.opensearch.analytics.planner.dag;
+
+import org.apache.calcite.rel.RelDistribution;
+import org.opensearch.analytics.planner.rel.ShuffleImpl;
+
+import java.util.List;
+
+/**
+ * Exchange metadata extracted from exchange RelNodes during DAG construction.
+ * Describes how a child stage delivers data to its parent stage.
+ *
+ * @param distributionType distribution type from the exchange operator's trait
+ * @param shuffleImpl shuffle implementation (null for SINGLETON)
+ * @param partitionKeyIndices field indices for hash/range partitioning (empty for SINGLETON)
+ * @opensearch.internal
+ */
+public record ExchangeInfo(
+ RelDistribution.Type distributionType,
+ ShuffleImpl shuffleImpl,
+ List partitionKeyIndices
+) {
+ public ExchangeInfo {
+ partitionKeyIndices = List.copyOf(partitionKeyIndices);
+ }
+
+ public boolean isShuffle() {
+ return shuffleImpl != null;
+ }
+}
diff --git a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/dag/QueryDAG.java b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/dag/QueryDAG.java
new file mode 100644
index 0000000000000..b2d4f17982d38
--- /dev/null
+++ b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/dag/QueryDAG.java
@@ -0,0 +1,61 @@
+/*
+ * SPDX-License-Identifier: Apache-2.0
+ *
+ * The OpenSearch Contributors require contributions made to
+ * this file be licensed under the Apache-2.0 license or a
+ * compatible open source license.
+ */
+
+package org.opensearch.analytics.planner.dag;
+
+import org.apache.calcite.plan.RelOptUtil;
+
+/**
+ * Root of the query execution DAG. Recursive tree of {@link Stage}s.
+ *
+ * @param queryId unique identifier for this query. Currently a random UUID.
+ * Future: accept a user-provided ID or generate a cluster-wide
+ * unique ID (e.g., node-local counter + nodeId prefix).
+ * @param rootStage the coordinator/root stage
+ * @opensearch.internal
+ */
+public record QueryDAG(String queryId, Stage rootStage) {
+
+ /**
+ * Returns a human-readable representation of the entire DAG showing
+ * every stage with its fragment indented.
+ */
+ @Override
+ public String toString() {
+ StringBuilder builder = new StringBuilder();
+ builder.append("QueryDAG(queryId=").append(queryId).append(")\n");
+ appendStage(builder, rootStage, 0);
+ return builder.toString();
+ }
+
+ private static void appendStage(StringBuilder builder, Stage stage, int depth) {
+ String indent = " ".repeat(depth);
+ builder.append(indent).append("Stage ").append(stage.getStageId());
+ if (stage.getExchangeInfo() != null) {
+ builder.append(" [exchange=").append(stage.getExchangeInfo()).append("]");
+ } else {
+ builder.append(" [root]");
+ }
+ builder.append("\n");
+
+ if (stage.getFragment() != null) {
+ String fragmentStr = RelOptUtil.toString(stage.getFragment());
+ for (String line : fragmentStr.split("\n")) {
+ if (!line.isEmpty()) {
+ builder.append(indent).append(" ").append(line).append("\n");
+ }
+ }
+ } else {
+ builder.append(indent).append(" \n");
+ }
+
+ for (Stage child : stage.getChildStages()) {
+ appendStage(builder, child, depth + 1);
+ }
+ }
+}
diff --git a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/dag/Stage.java b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/dag/Stage.java
new file mode 100644
index 0000000000000..255c173f5e416
--- /dev/null
+++ b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/dag/Stage.java
@@ -0,0 +1,57 @@
+/*
+ * SPDX-License-Identifier: Apache-2.0
+ *
+ * The OpenSearch Contributors require contributions made to
+ * this file be licensed under the Apache-2.0 license or a
+ * compatible open source license.
+ */
+
+package org.opensearch.analytics.planner.dag;
+
+import org.apache.calcite.rel.RelNode;
+
+import java.util.List;
+
+/**
+ * A stage in the query DAG. Each stage holds a plan fragment (the marked RelNode
+ * subtree between exchange boundaries) and references to child stages.
+ *
+ * Fragment may be null for a pure gather stage (coordinator just accumulates
+ * Arrow batches from child stages).
+ *
+ * @opensearch.internal
+ */
+public class Stage {
+
+ private final int stageId;
+ private final RelNode fragment;
+ private final List childStages;
+ private final ExchangeInfo exchangeInfo;
+
+ // TODO: add List planAlternatives — populated during plan forking phase
+
+ public Stage(int stageId, RelNode fragment, List childStages, ExchangeInfo exchangeInfo) {
+ this.stageId = stageId;
+ this.fragment = fragment;
+ this.childStages = List.copyOf(childStages);
+ this.exchangeInfo = exchangeInfo;
+ }
+
+ public int getStageId() {
+ return stageId;
+ }
+
+ /** Marked plan fragment with annotations intact. Null for a pure gather stage. */
+ public RelNode getFragment() {
+ return fragment;
+ }
+
+ public List getChildStages() {
+ return childStages;
+ }
+
+ /** How this stage connects to its parent. Null for the root stage. */
+ public ExchangeInfo getExchangeInfo() {
+ return exchangeInfo;
+ }
+}
diff --git a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/rel/OpenSearchStageInputScan.java b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/rel/OpenSearchStageInputScan.java
new file mode 100644
index 0000000000000..7b1eb24056f84
--- /dev/null
+++ b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/rel/OpenSearchStageInputScan.java
@@ -0,0 +1,69 @@
+/*
+ * SPDX-License-Identifier: Apache-2.0
+ *
+ * The OpenSearch Contributors require contributions made to
+ * this file be licensed under the Apache-2.0 license or a
+ * compatible open source license.
+ */
+
+package org.opensearch.analytics.planner.rel;
+
+import org.apache.calcite.plan.RelOptCluster;
+import org.apache.calcite.plan.RelOptCost;
+import org.apache.calcite.plan.RelOptPlanner;
+import org.apache.calcite.plan.RelTraitSet;
+import org.apache.calcite.rel.AbstractRelNode;
+import org.apache.calcite.rel.RelNode;
+import org.apache.calcite.rel.RelWriter;
+import org.apache.calcite.rel.metadata.RelMetadataQuery;
+import org.apache.calcite.rel.type.RelDataType;
+
+import java.util.List;
+
+/**
+ * Leaf node representing data arriving from a child stage via exchange.
+ * Replaces the severed child subtree in the parent stage fragment during
+ * DAG construction.
+ *
+ * Carries the child stage ID and the row type of the child subtree's
+ * output. During fragment conversion (RelNode → Substrait/QueryBuilder),
+ * the backend registers a streaming source or shuffle source for this node.
+ *
+ * @opensearch.internal
+ */
+public class OpenSearchStageInputScan extends AbstractRelNode {
+
+ private final int childStageId;
+ private final RelDataType rowType;
+
+ public OpenSearchStageInputScan(RelOptCluster cluster, RelTraitSet traitSet,
+ int childStageId, RelDataType rowType) {
+ super(cluster, traitSet);
+ this.childStageId = childStageId;
+ this.rowType = rowType;
+ }
+
+ public int getChildStageId() {
+ return childStageId;
+ }
+
+ @Override
+ protected RelDataType deriveRowType() {
+ return rowType;
+ }
+
+ @Override
+ public RelNode copy(RelTraitSet traitSet, List inputs) {
+ return new OpenSearchStageInputScan(getCluster(), traitSet, childStageId, rowType);
+ }
+
+ @Override
+ public RelOptCost computeSelfCost(RelOptPlanner planner, RelMetadataQuery mq) {
+ return planner.getCostFactory().makeZeroCost();
+ }
+
+ @Override
+ public RelWriter explainTerms(RelWriter pw) {
+ return super.explainTerms(pw).item("childStageId", childStageId);
+ }
+}
diff --git a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/rules/OpenSearchProjectRule.java b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/rules/OpenSearchProjectRule.java
index 161291de1f85a..73317d87bc6a3 100644
--- a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/rules/OpenSearchProjectRule.java
+++ b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/planner/rules/OpenSearchProjectRule.java
@@ -23,6 +23,7 @@
import org.opensearch.analytics.planner.rel.OpenSearchRelNode;
import org.opensearch.analytics.spi.DelegationType;
import org.opensearch.analytics.spi.FieldType;
+import org.opensearch.analytics.spi.OperatorCapability;
import org.opensearch.analytics.spi.ScalarFunction;
import java.util.ArrayList;
@@ -190,7 +191,11 @@ private List computeProjectViableBackends(List annotatedExprs,
List delegationAcceptors = registry.delegationAcceptors(DelegationType.PROJECT);
List result = new ArrayList<>();
+ List projectCapable = registry.operatorBackends(OperatorCapability.PROJECT);
for (String candidateName : childViableBackends) {
+ if (!projectCapable.contains(candidateName)) {
+ continue;
+ }
boolean canHandleAll = true;
for (RexNode expr : annotatedExprs) {
if (!(expr instanceof AnnotatedProjectExpression annotation)) {
diff --git a/sandbox/plugins/analytics-engine/src/test/java/org/opensearch/analytics/exec/DefaultPlanExecutorTests.java b/sandbox/plugins/analytics-engine/src/test/java/org/opensearch/analytics/exec/DefaultPlanExecutorTests.java
index 14d9dd2a17ed6..459a7f3786abd 100644
--- a/sandbox/plugins/analytics-engine/src/test/java/org/opensearch/analytics/exec/DefaultPlanExecutorTests.java
+++ b/sandbox/plugins/analytics-engine/src/test/java/org/opensearch/analytics/exec/DefaultPlanExecutorTests.java
@@ -26,11 +26,14 @@
import org.opensearch.analytics.backend.ExecutionContext;
import org.opensearch.analytics.backend.SearchExecEngine;
import org.opensearch.analytics.spi.AnalyticsSearchBackendPlugin;
+import org.opensearch.analytics.spi.OperatorCapability;
import org.opensearch.cluster.ClusterState;
import org.opensearch.cluster.metadata.IndexMetadata;
+import org.opensearch.cluster.metadata.MappingMetadata;
import org.opensearch.cluster.metadata.Metadata;
import org.opensearch.cluster.service.ClusterService;
import org.opensearch.common.concurrent.GatedCloseable;
+import org.opensearch.common.settings.Settings;
import org.opensearch.core.index.Index;
import org.opensearch.index.IndexService;
import org.opensearch.index.engine.DataFormatAwareEngine;
@@ -123,6 +126,15 @@ public void testEndToEndExecuteWithMockBackend() throws IOException {
Index index = new Index("my_index", "uuid");
IndexMetadata indexMetadata = mock(IndexMetadata.class);
when(indexMetadata.getIndex()).thenReturn(index);
+ when(indexMetadata.getSettings()).thenReturn(
+ Settings.builder().put("index.composite.primary_data_format", "mock-columnar").build()
+ );
+ when(indexMetadata.getNumberOfShards()).thenReturn(1);
+ MappingMetadata mappingMetadata = mock(MappingMetadata.class);
+ when(mappingMetadata.sourceAsMap()).thenReturn(
+ Map.of("properties", Map.of("field_0", Map.of("type", "integer")))
+ );
+ when(indexMetadata.mapping()).thenReturn(mappingMetadata);
Metadata metadata = mock(Metadata.class);
when(metadata.index("my_index")).thenReturn(indexMetadata);
@@ -398,6 +410,16 @@ public String name() {
return "mock-backend";
}
+ @Override
+ public List getSupportedFormats() {
+ return List.of(format);
+ }
+
+ @Override
+ public Set supportedOperators() {
+ return Set.of(OperatorCapability.SCAN, OperatorCapability.FILTER);
+ }
+
@Override
public SearchExecEngine createSearchExecEngine(ExecutionContext ctx) {
Object reader = ctx.getReader().reader(format);
diff --git a/sandbox/plugins/analytics-engine/src/test/java/org/opensearch/analytics/planner/MockDataFusionBackend.java b/sandbox/plugins/analytics-engine/src/test/java/org/opensearch/analytics/planner/MockDataFusionBackend.java
index 8335338c14315..c7e36ef7789ec 100644
--- a/sandbox/plugins/analytics-engine/src/test/java/org/opensearch/analytics/planner/MockDataFusionBackend.java
+++ b/sandbox/plugins/analytics-engine/src/test/java/org/opensearch/analytics/planner/MockDataFusionBackend.java
@@ -90,7 +90,7 @@ public class MockDataFusionBackend implements AnalyticsSearchBackendPlugin {
@Override public String name() { return NAME; }
- @Override public SearchExecEngine searcher(ExecutionContext ctx) { return null; }
+ @Override public SearchExecEngine createSearchExecEngine(ExecutionContext ctx) { return null; }
@Override
public List getSupportedFormats() {
diff --git a/sandbox/plugins/analytics-engine/src/test/java/org/opensearch/analytics/planner/MockLuceneBackend.java b/sandbox/plugins/analytics-engine/src/test/java/org/opensearch/analytics/planner/MockLuceneBackend.java
index 4b34e36e93dec..7828ae9f46f59 100644
--- a/sandbox/plugins/analytics-engine/src/test/java/org/opensearch/analytics/planner/MockLuceneBackend.java
+++ b/sandbox/plugins/analytics-engine/src/test/java/org/opensearch/analytics/planner/MockLuceneBackend.java
@@ -92,7 +92,7 @@ public class MockLuceneBackend implements AnalyticsSearchBackendPlugin {
@Override public String name() { return NAME; }
- @Override public SearchExecEngine searcher(ExecutionContext ctx) { return null; }
+ @Override public SearchExecEngine createSearchExecEngine(ExecutionContext ctx) { return null; }
@Override
public List getSupportedFormats() {
diff --git a/server/src/main/java/org/opensearch/index/mapper/FieldValueFetcher.java b/server/src/main/java/org/opensearch/index/mapper/FieldValueFetcher.java
index aeacf235591ad..82df8a4b9d08f 100644
--- a/server/src/main/java/org/opensearch/index/mapper/FieldValueFetcher.java
+++ b/server/src/main/java/org/opensearch/index/mapper/FieldValueFetcher.java
@@ -40,7 +40,7 @@ protected FieldValueFetcher(String simpleName) {
* Converts the field value to required representation, should be overridden by field mappers as needed
* @param value - value to convert
*/
- Object convert(Object value) {
+ public Object convert(Object value) {
return value;
}