Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line numberDiff line numberDiff line change
Expand Up@@ -36,7 +36,6 @@
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.time.Instant;
import java.util.List;
import java.util.Properties;
import java.util.function.Consumer;
import org.apache.calcite.adapter.enumerable.EnumerableConvention;
Expand All@@ -59,7 +58,6 @@
import org.apache.calcite.plan.Convention;
import org.apache.calcite.plan.RelOptCluster;
import org.apache.calcite.plan.RelOptPlanner;
import org.apache.calcite.plan.RelOptRule;
import org.apache.calcite.plan.RelOptSchema;
import org.apache.calcite.plan.RelOptTable;
import org.apache.calcite.plan.RelOptTable.ViewExpander;
Expand DownExpand Up@@ -551,12 +549,26 @@ public RelNode visit(TableScan scan) {
}
}

/** Try to optimize the plan by using HepPlanner */
private static final List<RelOptRule> hepRuleList =

Copy link
Copy Markdown
Owner

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ryan-gh-bot just for my knowledge why are we removing it?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

That hepRuleList field only existed as the argument to the old addRuleCollection(hepRuleList) call. Once the builder switched to two sequential addRuleInstance(...) calls (which take single rules directly), nothing else referenced the list, so keeping it would just be dead code — and dropping it also lets the java.util.List and org.apache.calcite.plan.RelOptRule imports go away. The Javadoc that used to sit on the field moved up onto HEP_PROGRAM, which is now the thing whose ordering actually matters.

List.of(FilterMergeRule.Config.DEFAULT.toRule(), PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE);

/**
* Try to optimize the plan by using HepPlanner.
*
* <p>Rule order matters: {@link PPLSimplifyDedupRule#DEDUP_SIMPLIFY_RULE} must run to fixpoint
* before {@link FilterMergeRule}. The simplify rule's bottom operand only matches a pure {@code
* IS NOT NULL} (or AND-of-{@code IS NOT NULL}) bucket-non-null filter; if {@code FilterMergeRule}
* runs first when a user {@code WHERE} sits below the synthetic {@code IS NOT NULL} filter that
* PPL emits as part of {@code dedup}, the two adjacent filters are merged into a single filter
* whose condition includes the user predicate, the simplify rule's predicate fails, no {@link
* org.opensearch.sql.calcite.plan.rel.LogicalDedup} is produced, and dedup pushdown to the
* OpenSearch storage engine is silently disabled. Using separate {@code addRuleInstance} calls
* (rather than {@code addRuleCollection}) enforces deterministic ordering: dedup simplification
* fires first against the original adjacent-filter shape, then any remaining adjacent filters are
* merged.
*/
private static final HepProgram HEP_PROGRAM =
new HepProgramBuilder().addRuleCollection(hepRuleList).build();
new HepProgramBuilder()
.addRuleInstance(PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE)
.addRuleInstance(FilterMergeRule.Config.DEFAULT.toRule())
.build();

public static RelNode optimize(RelNode plan, CalcitePlanContext context) {
Util.discard(context);
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -45,6 +45,7 @@
import org.opensearch.sql.calcite.CalcitePlanContext;
import org.opensearch.sql.calcite.CalciteRelNodeVisitor;
import org.opensearch.sql.calcite.SysLimit;
import org.opensearch.sql.calcite.utils.CalciteToolsHelper;
import org.opensearch.sql.common.setting.Settings;
import org.opensearch.sql.datasource.DataSourceService;
import org.opensearch.sql.exception.ExpressionEvaluationException;
Expand DownExpand Up@@ -110,6 +111,20 @@ public RelNode getRelNode(String ppl) {
return root;
}

/**
* Get the root RelNode of the given PPL query after running the production HEP program from
* {@code CalciteToolsHelper}. Use this in regression tests that exercise rules registered in the
* production HEP program (e.g. {@code PPLSimplifyDedupRule}) — those rules need to see the raw
* planner output, not the post-FilterMerge form returned by {@link #getRelNode(String)}.
*/
public RelNode getRelNodeAfterCalciteHep(String ppl) {
CalcitePlanContext context = createBuilderContext();
Query query = (Query) plan(pplParser, ppl);
planTransformer.analyze(query.getPlan(), context);
RelNode root = context.relBuilder.build();
return CalciteToolsHelper.optimize(root, context);
}

private RelNode mergeAdjacentFilters(RelNode relNode) {
HepProgram program =
new HepProgramBuilder().addRuleInstance(FilterMergeRule.Config.DEFAULT.toRule()).build();
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -5,9 +5,17 @@

package org.opensearch.sql.ppl.calcite;

import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;

import org.apache.calcite.plan.hep.HepPlanner;
import org.apache.calcite.plan.hep.HepProgram;
import org.apache.calcite.plan.hep.HepProgramBuilder;
import org.apache.calcite.rel.RelNode;
import org.apache.calcite.test.CalciteAssert;
import org.junit.Test;
import org.opensearch.sql.calcite.plan.rel.LogicalDedup;
import org.opensearch.sql.calcite.plan.rule.PPLSimplifyDedupRule;

public class CalcitePPLDedupTest extends CalcitePPLAbstractTest {

Expand DownExpand Up@@ -353,4 +361,72 @@ public void testSortFieldProjectedAwayBeforeDedup() {
+ " LogicalTableScan(table=[[scott, EMP]])\n";
verifyLogical(root, expectedLogical);
}

/**
* Regression test for issue #7: when a user {@code where} sits below {@code dedup}, the HEP
* program in {@code CalciteToolsHelper} must still produce a {@link LogicalDedup}. Before the
* fix, both rules were registered via {@code addRuleCollection}, so {@code FilterMergeRule} could
* fire ahead of {@code PPLSimplifyDedupRule} and merge the user predicate into the
* bucket-non-null filter; the simplify rule's bottom operand then rejected the merged condition
* (it only accepts pure {@code IS NOT NULL}/AND-of-{@code IS NOT NULL}), no {@code LogicalDedup}
* was produced, and dedup pushdown to the OpenSearch storage engine was silently disabled. The
* fix is to register the two rules with separate {@code addRuleInstance} calls in the order
* simplify-dedup first (to fixpoint), then filter-merge.
*/
@Test
public void testWhereThenDedupProducesLogicalDedup() {
// Use a where predicate on a DIFFERENT column from the dedup column. With the same column,
// Calcite's RexSimplify can fold AND(IS_NOT_NULL(x), >(x, c)) down to >(x, c), masking the
// bug. The issue's reproducer (where on @timestamp, dedup on namespace) hits this exact
// shape.
String ppl = "source=EMP | where SAL > 1000 | dedup 1 DEPTNO | fields DEPTNO";
RelNode optimized = getRelNodeAfterCalciteHep(ppl);
String optimizedPlan = optimized.explain();
assertTrue(
"where + dedup must produce a LogicalDedup so OpenSearch DedupPushdownRule can match;"
+ " actual plan was:\n"
+ optimizedPlan,
optimizedPlan.contains("LogicalDedup"));
// The window-form leftover would indicate the simplify rule did not fire — assert it is gone.
assertFalse(
"ROW_NUMBER window must be consumed by PPLSimplifyDedupRule when where + dedup are"
+ " combined; actual plan was:\n"
+ optimizedPlan,
optimizedPlan.contains("ROW_NUMBER"));
}

/**
* Adversarial regression test: simulates the pathological order described in issue #7 by forcing
* FilterMergeRule to run to fixpoint BEFORE PPLSimplifyDedupRule. This documents the failure mode
* the fix in {@code CalciteToolsHelper} prevents — once the bucket-non-null filter has been
* merged with the user {@code WHERE}, {@code mayBeFilterFromBucketNonNull} can never accept the
* combined condition, so {@code PPLSimplifyDedupRule} is permanently unable to produce a {@code
* LogicalDedup}. The production fix enforces order at the program level (sequential {@code
* addRuleInstance} calls), making this hazard unreachable.
*/
@Test
public void testFilterMergeBeforeSimplifyDedupBreaksPattern() {
String ppl = "source=EMP | where SAL > 1000 | dedup 1 DEPTNO | fields DEPTNO";
// getRelNode already runs FilterMergeRule on the raw plan, simulating the pathological
// schedule where FilterMergeRule fires before PPLSimplifyDedupRule.
RelNode mergedFirst = getRelNode(ppl);
HepProgram simplifyOnly =
new HepProgramBuilder().addRuleInstance(PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE).build();
HepPlanner planner = new HepPlanner(simplifyOnly);
planner.setRoot(mergedFirst);
RelNode result = planner.findBestExp();
String plan = result.explain();
assertFalse(
"If FilterMergeRule runs before PPLSimplifyDedupRule, the simplify rule must NOT recover"
+ " — the merged AND(IS_NOT_NULL, user_predicate) filter fails the bucket-non-null"
+ " predicate. This documents why the production HEP program enforces ordering via"
+ " separate addRuleInstance calls (PPLSimplifyDedupRule first, then FilterMergeRule)."
+ " Actual plan was:\n"
+ plan,
plan.contains("LogicalDedup"));
assertTrue(
"Plan should still contain ROW_NUMBER window form when simplify fails. Actual plan was:\n"
+ plan,
plan.contains("ROW_NUMBER"));
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line numberDiff line numberDiff line change
Expand Up@@ -36,7 +36,6 @@
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.time.Instant;
import java.util.List;
import java.util.Properties;
import java.util.function.Consumer;
import org.apache.calcite.adapter.enumerable.EnumerableConvention;
Expand All@@ -59,7 +58,6 @@
import org.apache.calcite.plan.Convention;
import org.apache.calcite.plan.RelOptCluster;
import org.apache.calcite.plan.RelOptPlanner;
import org.apache.calcite.plan.RelOptRule;
import org.apache.calcite.plan.RelOptSchema;
import org.apache.calcite.plan.RelOptTable;
import org.apache.calcite.plan.RelOptTable.ViewExpander;
Expand DownExpand Up@@ -551,12 +549,26 @@ public RelNode visit(TableScan scan) {
}
}

/** Try to optimize the plan by using HepPlanner */
private static final List<RelOptRule> hepRuleList =

Copy link
Copy Markdown
Owner

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ryan-gh-bot just for my knowledge why are we removing it?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

That hepRuleList field only existed as the argument to the old addRuleCollection(hepRuleList) call. Once the builder switched to two sequential addRuleInstance(...) calls (which take single rules directly), nothing else referenced the list, so keeping it would just be dead code — and dropping it also lets the java.util.List and org.apache.calcite.plan.RelOptRule imports go away. The Javadoc that used to sit on the field moved up onto HEP_PROGRAM, which is now the thing whose ordering actually matters.

List.of(FilterMergeRule.Config.DEFAULT.toRule(), PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE);

/**
* Try to optimize the plan by using HepPlanner.
*
* <p>Rule order matters: {@link PPLSimplifyDedupRule#DEDUP_SIMPLIFY_RULE} must run to fixpoint
* before {@link FilterMergeRule}. The simplify rule's bottom operand only matches a pure {@code
* IS NOT NULL} (or AND-of-{@code IS NOT NULL}) bucket-non-null filter; if {@code FilterMergeRule}
* runs first when a user {@code WHERE} sits below the synthetic {@code IS NOT NULL} filter that
* PPL emits as part of {@code dedup}, the two adjacent filters are merged into a single filter
* whose condition includes the user predicate, the simplify rule's predicate fails, no {@link
* org.opensearch.sql.calcite.plan.rel.LogicalDedup} is produced, and dedup pushdown to the
* OpenSearch storage engine is silently disabled. Using separate {@code addRuleInstance} calls
* (rather than {@code addRuleCollection}) enforces deterministic ordering: dedup simplification
* fires first against the original adjacent-filter shape, then any remaining adjacent filters are
* merged.
*/
private static final HepProgram HEP_PROGRAM =
new HepProgramBuilder().addRuleCollection(hepRuleList).build();
new HepProgramBuilder()
.addRuleInstance(PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE)
.addRuleInstance(FilterMergeRule.Config.DEFAULT.toRule())
.build();

public static RelNode optimize(RelNode plan, CalcitePlanContext context) {
Util.discard(context);
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -45,6 +45,7 @@
import org.opensearch.sql.calcite.CalcitePlanContext;
import org.opensearch.sql.calcite.CalciteRelNodeVisitor;
import org.opensearch.sql.calcite.SysLimit;
import org.opensearch.sql.calcite.utils.CalciteToolsHelper;
import org.opensearch.sql.common.setting.Settings;
import org.opensearch.sql.datasource.DataSourceService;
import org.opensearch.sql.exception.ExpressionEvaluationException;
Expand DownExpand Up@@ -110,6 +111,20 @@ public RelNode getRelNode(String ppl) {
return root;
}

/**
* Get the root RelNode of the given PPL query after running the production HEP program from
* {@code CalciteToolsHelper}. Use this in regression tests that exercise rules registered in the
* production HEP program (e.g. {@code PPLSimplifyDedupRule}) — those rules need to see the raw
* planner output, not the post-FilterMerge form returned by {@link #getRelNode(String)}.
*/
public RelNode getRelNodeAfterCalciteHep(String ppl) {
CalcitePlanContext context = createBuilderContext();
Query query = (Query) plan(pplParser, ppl);
planTransformer.analyze(query.getPlan(), context);
RelNode root = context.relBuilder.build();
return CalciteToolsHelper.optimize(root, context);
}

private RelNode mergeAdjacentFilters(RelNode relNode) {
HepProgram program =
new HepProgramBuilder().addRuleInstance(FilterMergeRule.Config.DEFAULT.toRule()).build();
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -5,9 +5,17 @@

package org.opensearch.sql.ppl.calcite;

import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;

import org.apache.calcite.plan.hep.HepPlanner;
import org.apache.calcite.plan.hep.HepProgram;
import org.apache.calcite.plan.hep.HepProgramBuilder;
import org.apache.calcite.rel.RelNode;
import org.apache.calcite.test.CalciteAssert;
import org.junit.Test;
import org.opensearch.sql.calcite.plan.rel.LogicalDedup;
import org.opensearch.sql.calcite.plan.rule.PPLSimplifyDedupRule;

public class CalcitePPLDedupTest extends CalcitePPLAbstractTest {

Expand DownExpand Up@@ -353,4 +361,72 @@ public void testSortFieldProjectedAwayBeforeDedup() {
+ " LogicalTableScan(table=[[scott, EMP]])\n";
verifyLogical(root, expectedLogical);
}

/**
* Regression test for issue #7: when a user {@code where} sits below {@code dedup}, the HEP
* program in {@code CalciteToolsHelper} must still produce a {@link LogicalDedup}. Before the
* fix, both rules were registered via {@code addRuleCollection}, so {@code FilterMergeRule} could
* fire ahead of {@code PPLSimplifyDedupRule} and merge the user predicate into the
* bucket-non-null filter; the simplify rule's bottom operand then rejected the merged condition
* (it only accepts pure {@code IS NOT NULL}/AND-of-{@code IS NOT NULL}), no {@code LogicalDedup}
* was produced, and dedup pushdown to the OpenSearch storage engine was silently disabled. The
* fix is to register the two rules with separate {@code addRuleInstance} calls in the order
* simplify-dedup first (to fixpoint), then filter-merge.
*/
@Test
public void testWhereThenDedupProducesLogicalDedup() {
// Use a where predicate on a DIFFERENT column from the dedup column. With the same column,
// Calcite's RexSimplify can fold AND(IS_NOT_NULL(x), >(x, c)) down to >(x, c), masking the
// bug. The issue's reproducer (where on @timestamp, dedup on namespace) hits this exact
// shape.
String ppl = "source=EMP | where SAL > 1000 | dedup 1 DEPTNO | fields DEPTNO";
RelNode optimized = getRelNodeAfterCalciteHep(ppl);
String optimizedPlan = optimized.explain();
assertTrue(
"where + dedup must produce a LogicalDedup so OpenSearch DedupPushdownRule can match;"
+ " actual plan was:\n"
+ optimizedPlan,
optimizedPlan.contains("LogicalDedup"));
// The window-form leftover would indicate the simplify rule did not fire — assert it is gone.
assertFalse(
"ROW_NUMBER window must be consumed by PPLSimplifyDedupRule when where + dedup are"
+ " combined; actual plan was:\n"
+ optimizedPlan,
optimizedPlan.contains("ROW_NUMBER"));
}

/**
* Adversarial regression test: simulates the pathological order described in issue #7 by forcing
* FilterMergeRule to run to fixpoint BEFORE PPLSimplifyDedupRule. This documents the failure mode
* the fix in {@code CalciteToolsHelper} prevents — once the bucket-non-null filter has been
* merged with the user {@code WHERE}, {@code mayBeFilterFromBucketNonNull} can never accept the
* combined condition, so {@code PPLSimplifyDedupRule} is permanently unable to produce a {@code
* LogicalDedup}. The production fix enforces order at the program level (sequential {@code
* addRuleInstance} calls), making this hazard unreachable.
*/
@Test
public void testFilterMergeBeforeSimplifyDedupBreaksPattern() {
String ppl = "source=EMP | where SAL > 1000 | dedup 1 DEPTNO | fields DEPTNO";
// getRelNode already runs FilterMergeRule on the raw plan, simulating the pathological
// schedule where FilterMergeRule fires before PPLSimplifyDedupRule.
RelNode mergedFirst = getRelNode(ppl);
HepProgram simplifyOnly =
new HepProgramBuilder().addRuleInstance(PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE).build();
HepPlanner planner = new HepPlanner(simplifyOnly);
planner.setRoot(mergedFirst);
RelNode result = planner.findBestExp();
String plan = result.explain();
assertFalse(
"If FilterMergeRule runs before PPLSimplifyDedupRule, the simplify rule must NOT recover"
+ " — the merged AND(IS_NOT_NULL, user_predicate) filter fails the bucket-non-null"
+ " predicate. This documents why the production HEP program enforces ordering via"
+ " separate addRuleInstance calls (PPLSimplifyDedupRule first, then FilterMergeRule)."
+ " Actual plan was:\n"
+ plan,
plan.contains("LogicalDedup"));
assertTrue(
"Plan should still contain ROW_NUMBER window form when simplify fails. Actual plan was:\n"
+ plan,
plan.contains("ROW_NUMBER"));
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line numberDiff line numberDiff line change
Expand Up@@ -36,7 +36,6 @@
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.time.Instant;
import java.util.List;
import java.util.Properties;
import java.util.function.Consumer;
import org.apache.calcite.adapter.enumerable.EnumerableConvention;
Expand All@@ -59,7 +58,6 @@
import org.apache.calcite.plan.Convention;
import org.apache.calcite.plan.RelOptCluster;
import org.apache.calcite.plan.RelOptPlanner;
import org.apache.calcite.plan.RelOptRule;
import org.apache.calcite.plan.RelOptSchema;
import org.apache.calcite.plan.RelOptTable;
import org.apache.calcite.plan.RelOptTable.ViewExpander;
Expand DownExpand Up@@ -551,12 +549,26 @@ public RelNode visit(TableScan scan) {
}
}

/** Try to optimize the plan by using HepPlanner */
private static final List<RelOptRule> hepRuleList =

Copy link
Copy Markdown
Owner

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ryan-gh-bot just for my knowledge why are we removing it?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

That hepRuleList field only existed as the argument to the old addRuleCollection(hepRuleList) call. Once the builder switched to two sequential addRuleInstance(...) calls (which take single rules directly), nothing else referenced the list, so keeping it would just be dead code — and dropping it also lets the java.util.List and org.apache.calcite.plan.RelOptRule imports go away. The Javadoc that used to sit on the field moved up onto HEP_PROGRAM, which is now the thing whose ordering actually matters.

List.of(FilterMergeRule.Config.DEFAULT.toRule(), PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE);

/**
* Try to optimize the plan by using HepPlanner.
*
* <p>Rule order matters: {@link PPLSimplifyDedupRule#DEDUP_SIMPLIFY_RULE} must run to fixpoint
* before {@link FilterMergeRule}. The simplify rule's bottom operand only matches a pure {@code
* IS NOT NULL} (or AND-of-{@code IS NOT NULL}) bucket-non-null filter; if {@code FilterMergeRule}
* runs first when a user {@code WHERE} sits below the synthetic {@code IS NOT NULL} filter that
* PPL emits as part of {@code dedup}, the two adjacent filters are merged into a single filter
* whose condition includes the user predicate, the simplify rule's predicate fails, no {@link
* org.opensearch.sql.calcite.plan.rel.LogicalDedup} is produced, and dedup pushdown to the
* OpenSearch storage engine is silently disabled. Using separate {@code addRuleInstance} calls
* (rather than {@code addRuleCollection}) enforces deterministic ordering: dedup simplification
* fires first against the original adjacent-filter shape, then any remaining adjacent filters are
* merged.
*/
private static final HepProgram HEP_PROGRAM =
new HepProgramBuilder().addRuleCollection(hepRuleList).build();
new HepProgramBuilder()
.addRuleInstance(PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE)
.addRuleInstance(FilterMergeRule.Config.DEFAULT.toRule())
.build();

public static RelNode optimize(RelNode plan, CalcitePlanContext context) {
Util.discard(context);
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -45,6 +45,7 @@
import org.opensearch.sql.calcite.CalcitePlanContext;
import org.opensearch.sql.calcite.CalciteRelNodeVisitor;
import org.opensearch.sql.calcite.SysLimit;
import org.opensearch.sql.calcite.utils.CalciteToolsHelper;
import org.opensearch.sql.common.setting.Settings;
import org.opensearch.sql.datasource.DataSourceService;
import org.opensearch.sql.exception.ExpressionEvaluationException;
Expand DownExpand Up@@ -110,6 +111,20 @@ public RelNode getRelNode(String ppl) {
return root;
}

/**
* Get the root RelNode of the given PPL query after running the production HEP program from
* {@code CalciteToolsHelper}. Use this in regression tests that exercise rules registered in the
* production HEP program (e.g. {@code PPLSimplifyDedupRule}) — those rules need to see the raw
* planner output, not the post-FilterMerge form returned by {@link #getRelNode(String)}.
*/
public RelNode getRelNodeAfterCalciteHep(String ppl) {
CalcitePlanContext context = createBuilderContext();
Query query = (Query) plan(pplParser, ppl);
planTransformer.analyze(query.getPlan(), context);
RelNode root = context.relBuilder.build();
return CalciteToolsHelper.optimize(root, context);
}

private RelNode mergeAdjacentFilters(RelNode relNode) {
HepProgram program =
new HepProgramBuilder().addRuleInstance(FilterMergeRule.Config.DEFAULT.toRule()).build();
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -5,9 +5,17 @@

package org.opensearch.sql.ppl.calcite;

import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;

import org.apache.calcite.plan.hep.HepPlanner;
import org.apache.calcite.plan.hep.HepProgram;
import org.apache.calcite.plan.hep.HepProgramBuilder;
import org.apache.calcite.rel.RelNode;
import org.apache.calcite.test.CalciteAssert;
import org.junit.Test;
import org.opensearch.sql.calcite.plan.rel.LogicalDedup;
import org.opensearch.sql.calcite.plan.rule.PPLSimplifyDedupRule;

public class CalcitePPLDedupTest extends CalcitePPLAbstractTest {

Expand DownExpand Up@@ -353,4 +361,72 @@ public void testSortFieldProjectedAwayBeforeDedup() {
+ " LogicalTableScan(table=[[scott, EMP]])\n";
verifyLogical(root, expectedLogical);
}

/**
* Regression test for issue #7: when a user {@code where} sits below {@code dedup}, the HEP
* program in {@code CalciteToolsHelper} must still produce a {@link LogicalDedup}. Before the
* fix, both rules were registered via {@code addRuleCollection}, so {@code FilterMergeRule} could
* fire ahead of {@code PPLSimplifyDedupRule} and merge the user predicate into the
* bucket-non-null filter; the simplify rule's bottom operand then rejected the merged condition
* (it only accepts pure {@code IS NOT NULL}/AND-of-{@code IS NOT NULL}), no {@code LogicalDedup}
* was produced, and dedup pushdown to the OpenSearch storage engine was silently disabled. The
* fix is to register the two rules with separate {@code addRuleInstance} calls in the order
* simplify-dedup first (to fixpoint), then filter-merge.
*/
@Test
public void testWhereThenDedupProducesLogicalDedup() {
// Use a where predicate on a DIFFERENT column from the dedup column. With the same column,
// Calcite's RexSimplify can fold AND(IS_NOT_NULL(x), >(x, c)) down to >(x, c), masking the
// bug. The issue's reproducer (where on @timestamp, dedup on namespace) hits this exact
// shape.
String ppl = "source=EMP | where SAL > 1000 | dedup 1 DEPTNO | fields DEPTNO";
RelNode optimized = getRelNodeAfterCalciteHep(ppl);
String optimizedPlan = optimized.explain();
assertTrue(
"where + dedup must produce a LogicalDedup so OpenSearch DedupPushdownRule can match;"
+ " actual plan was:\n"
+ optimizedPlan,
optimizedPlan.contains("LogicalDedup"));
// The window-form leftover would indicate the simplify rule did not fire — assert it is gone.
assertFalse(
"ROW_NUMBER window must be consumed by PPLSimplifyDedupRule when where + dedup are"
+ " combined; actual plan was:\n"
+ optimizedPlan,
optimizedPlan.contains("ROW_NUMBER"));
}

/**
* Adversarial regression test: simulates the pathological order described in issue #7 by forcing
* FilterMergeRule to run to fixpoint BEFORE PPLSimplifyDedupRule. This documents the failure mode
* the fix in {@code CalciteToolsHelper} prevents — once the bucket-non-null filter has been
* merged with the user {@code WHERE}, {@code mayBeFilterFromBucketNonNull} can never accept the
* combined condition, so {@code PPLSimplifyDedupRule} is permanently unable to produce a {@code
* LogicalDedup}. The production fix enforces order at the program level (sequential {@code
* addRuleInstance} calls), making this hazard unreachable.
*/
@Test
public void testFilterMergeBeforeSimplifyDedupBreaksPattern() {
String ppl = "source=EMP | where SAL > 1000 | dedup 1 DEPTNO | fields DEPTNO";
// getRelNode already runs FilterMergeRule on the raw plan, simulating the pathological
// schedule where FilterMergeRule fires before PPLSimplifyDedupRule.
RelNode mergedFirst = getRelNode(ppl);
HepProgram simplifyOnly =
new HepProgramBuilder().addRuleInstance(PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE).build();
HepPlanner planner = new HepPlanner(simplifyOnly);
planner.setRoot(mergedFirst);
RelNode result = planner.findBestExp();
String plan = result.explain();
assertFalse(
"If FilterMergeRule runs before PPLSimplifyDedupRule, the simplify rule must NOT recover"
+ " — the merged AND(IS_NOT_NULL, user_predicate) filter fails the bucket-non-null"
+ " predicate. This documents why the production HEP program enforces ordering via"
+ " separate addRuleInstance calls (PPLSimplifyDedupRule first, then FilterMergeRule)."
+ " Actual plan was:\n"
+ plan,
plan.contains("LogicalDedup"));
assertTrue(
"Plan should still contain ROW_NUMBER window form when simplify fails. Actual plan was:\n"
+ plan,
plan.contains("ROW_NUMBER"));
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line numberDiff line numberDiff line change
Expand Up@@ -36,7 +36,6 @@
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.time.Instant;
import java.util.List;
import java.util.Properties;
import java.util.function.Consumer;
import org.apache.calcite.adapter.enumerable.EnumerableConvention;
Expand All@@ -59,7 +58,6 @@
import org.apache.calcite.plan.Convention;
import org.apache.calcite.plan.RelOptCluster;
import org.apache.calcite.plan.RelOptPlanner;
import org.apache.calcite.plan.RelOptRule;
import org.apache.calcite.plan.RelOptSchema;
import org.apache.calcite.plan.RelOptTable;
import org.apache.calcite.plan.RelOptTable.ViewExpander;
Expand DownExpand Up@@ -551,12 +549,26 @@ public RelNode visit(TableScan scan) {
}
}

/** Try to optimize the plan by using HepPlanner */
private static final List<RelOptRule> hepRuleList =

Copy link
Copy Markdown
Owner

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ryan-gh-bot just for my knowledge why are we removing it?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

That hepRuleList field only existed as the argument to the old addRuleCollection(hepRuleList) call. Once the builder switched to two sequential addRuleInstance(...) calls (which take single rules directly), nothing else referenced the list, so keeping it would just be dead code — and dropping it also lets the java.util.List and org.apache.calcite.plan.RelOptRule imports go away. The Javadoc that used to sit on the field moved up onto HEP_PROGRAM, which is now the thing whose ordering actually matters.

List.of(FilterMergeRule.Config.DEFAULT.toRule(), PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE);

/**
* Try to optimize the plan by using HepPlanner.
*
* <p>Rule order matters: {@link PPLSimplifyDedupRule#DEDUP_SIMPLIFY_RULE} must run to fixpoint
* before {@link FilterMergeRule}. The simplify rule's bottom operand only matches a pure {@code
* IS NOT NULL} (or AND-of-{@code IS NOT NULL}) bucket-non-null filter; if {@code FilterMergeRule}
* runs first when a user {@code WHERE} sits below the synthetic {@code IS NOT NULL} filter that
* PPL emits as part of {@code dedup}, the two adjacent filters are merged into a single filter
* whose condition includes the user predicate, the simplify rule's predicate fails, no {@link
* org.opensearch.sql.calcite.plan.rel.LogicalDedup} is produced, and dedup pushdown to the
* OpenSearch storage engine is silently disabled. Using separate {@code addRuleInstance} calls
* (rather than {@code addRuleCollection}) enforces deterministic ordering: dedup simplification
* fires first against the original adjacent-filter shape, then any remaining adjacent filters are
* merged.
*/
private static final HepProgram HEP_PROGRAM =
new HepProgramBuilder().addRuleCollection(hepRuleList).build();
new HepProgramBuilder()
.addRuleInstance(PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE)
.addRuleInstance(FilterMergeRule.Config.DEFAULT.toRule())
.build();

public static RelNode optimize(RelNode plan, CalcitePlanContext context) {
Util.discard(context);
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -45,6 +45,7 @@
import org.opensearch.sql.calcite.CalcitePlanContext;
import org.opensearch.sql.calcite.CalciteRelNodeVisitor;
import org.opensearch.sql.calcite.SysLimit;
import org.opensearch.sql.calcite.utils.CalciteToolsHelper;
import org.opensearch.sql.common.setting.Settings;
import org.opensearch.sql.datasource.DataSourceService;
import org.opensearch.sql.exception.ExpressionEvaluationException;
Expand DownExpand Up@@ -110,6 +111,20 @@ public RelNode getRelNode(String ppl) {
return root;
}

/**
* Get the root RelNode of the given PPL query after running the production HEP program from
* {@code CalciteToolsHelper}. Use this in regression tests that exercise rules registered in the
* production HEP program (e.g. {@code PPLSimplifyDedupRule}) — those rules need to see the raw
* planner output, not the post-FilterMerge form returned by {@link #getRelNode(String)}.
*/
public RelNode getRelNodeAfterCalciteHep(String ppl) {
CalcitePlanContext context = createBuilderContext();
Query query = (Query) plan(pplParser, ppl);
planTransformer.analyze(query.getPlan(), context);
RelNode root = context.relBuilder.build();
return CalciteToolsHelper.optimize(root, context);
}

private RelNode mergeAdjacentFilters(RelNode relNode) {
HepProgram program =
new HepProgramBuilder().addRuleInstance(FilterMergeRule.Config.DEFAULT.toRule()).build();
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -5,9 +5,17 @@

package org.opensearch.sql.ppl.calcite;

import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;

import org.apache.calcite.plan.hep.HepPlanner;
import org.apache.calcite.plan.hep.HepProgram;
import org.apache.calcite.plan.hep.HepProgramBuilder;
import org.apache.calcite.rel.RelNode;
import org.apache.calcite.test.CalciteAssert;
import org.junit.Test;
import org.opensearch.sql.calcite.plan.rel.LogicalDedup;
import org.opensearch.sql.calcite.plan.rule.PPLSimplifyDedupRule;

public class CalcitePPLDedupTest extends CalcitePPLAbstractTest {

Expand DownExpand Up@@ -353,4 +361,72 @@ public void testSortFieldProjectedAwayBeforeDedup() {
+ " LogicalTableScan(table=[[scott, EMP]])\n";
verifyLogical(root, expectedLogical);
}

/**
* Regression test for issue #7: when a user {@code where} sits below {@code dedup}, the HEP
* program in {@code CalciteToolsHelper} must still produce a {@link LogicalDedup}. Before the
* fix, both rules were registered via {@code addRuleCollection}, so {@code FilterMergeRule} could
* fire ahead of {@code PPLSimplifyDedupRule} and merge the user predicate into the
* bucket-non-null filter; the simplify rule's bottom operand then rejected the merged condition
* (it only accepts pure {@code IS NOT NULL}/AND-of-{@code IS NOT NULL}), no {@code LogicalDedup}
* was produced, and dedup pushdown to the OpenSearch storage engine was silently disabled. The
* fix is to register the two rules with separate {@code addRuleInstance} calls in the order
* simplify-dedup first (to fixpoint), then filter-merge.
*/
@Test
public void testWhereThenDedupProducesLogicalDedup() {
// Use a where predicate on a DIFFERENT column from the dedup column. With the same column,
// Calcite's RexSimplify can fold AND(IS_NOT_NULL(x), >(x, c)) down to >(x, c), masking the
// bug. The issue's reproducer (where on @timestamp, dedup on namespace) hits this exact
// shape.
String ppl = "source=EMP | where SAL > 1000 | dedup 1 DEPTNO | fields DEPTNO";
RelNode optimized = getRelNodeAfterCalciteHep(ppl);
String optimizedPlan = optimized.explain();
assertTrue(
"where + dedup must produce a LogicalDedup so OpenSearch DedupPushdownRule can match;"
+ " actual plan was:\n"
+ optimizedPlan,
optimizedPlan.contains("LogicalDedup"));
// The window-form leftover would indicate the simplify rule did not fire — assert it is gone.
assertFalse(
"ROW_NUMBER window must be consumed by PPLSimplifyDedupRule when where + dedup are"
+ " combined; actual plan was:\n"
+ optimizedPlan,
optimizedPlan.contains("ROW_NUMBER"));
}

/**
* Adversarial regression test: simulates the pathological order described in issue #7 by forcing
* FilterMergeRule to run to fixpoint BEFORE PPLSimplifyDedupRule. This documents the failure mode
* the fix in {@code CalciteToolsHelper} prevents — once the bucket-non-null filter has been
* merged with the user {@code WHERE}, {@code mayBeFilterFromBucketNonNull} can never accept the
* combined condition, so {@code PPLSimplifyDedupRule} is permanently unable to produce a {@code
* LogicalDedup}. The production fix enforces order at the program level (sequential {@code
* addRuleInstance} calls), making this hazard unreachable.
*/
@Test
public void testFilterMergeBeforeSimplifyDedupBreaksPattern() {
String ppl = "source=EMP | where SAL > 1000 | dedup 1 DEPTNO | fields DEPTNO";
// getRelNode already runs FilterMergeRule on the raw plan, simulating the pathological
// schedule where FilterMergeRule fires before PPLSimplifyDedupRule.
RelNode mergedFirst = getRelNode(ppl);
HepProgram simplifyOnly =
new HepProgramBuilder().addRuleInstance(PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE).build();
HepPlanner planner = new HepPlanner(simplifyOnly);
planner.setRoot(mergedFirst);
RelNode result = planner.findBestExp();
String plan = result.explain();
assertFalse(
"If FilterMergeRule runs before PPLSimplifyDedupRule, the simplify rule must NOT recover"
+ " — the merged AND(IS_NOT_NULL, user_predicate) filter fails the bucket-non-null"
+ " predicate. This documents why the production HEP program enforces ordering via"
+ " separate addRuleInstance calls (PPLSimplifyDedupRule first, then FilterMergeRule)."
+ " Actual plan was:\n"
+ plan,
plan.contains("LogicalDedup"));
assertTrue(
"Plan should still contain ROW_NUMBER window form when simplify fails. Actual plan was:\n"
+ plan,
plan.contains("ROW_NUMBER"));
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line numberDiff line numberDiff line change
Expand Up@@ -36,7 +36,6 @@
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.time.Instant;
import java.util.List;
import java.util.Properties;
import java.util.function.Consumer;
import org.apache.calcite.adapter.enumerable.EnumerableConvention;
Expand All@@ -59,7 +58,6 @@
import org.apache.calcite.plan.Convention;
import org.apache.calcite.plan.RelOptCluster;
import org.apache.calcite.plan.RelOptPlanner;
import org.apache.calcite.plan.RelOptRule;
import org.apache.calcite.plan.RelOptSchema;
import org.apache.calcite.plan.RelOptTable;
import org.apache.calcite.plan.RelOptTable.ViewExpander;
Expand DownExpand Up@@ -551,12 +549,26 @@ public RelNode visit(TableScan scan) {
}
}

/** Try to optimize the plan by using HepPlanner */
private static final List<RelOptRule> hepRuleList =

Copy link
Copy Markdown
Owner

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ryan-gh-bot just for my knowledge why are we removing it?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

That hepRuleList field only existed as the argument to the old addRuleCollection(hepRuleList) call. Once the builder switched to two sequential addRuleInstance(...) calls (which take single rules directly), nothing else referenced the list, so keeping it would just be dead code — and dropping it also lets the java.util.List and org.apache.calcite.plan.RelOptRule imports go away. The Javadoc that used to sit on the field moved up onto HEP_PROGRAM, which is now the thing whose ordering actually matters.

List.of(FilterMergeRule.Config.DEFAULT.toRule(), PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE);

/**
* Try to optimize the plan by using HepPlanner.
*
* <p>Rule order matters: {@link PPLSimplifyDedupRule#DEDUP_SIMPLIFY_RULE} must run to fixpoint
* before {@link FilterMergeRule}. The simplify rule's bottom operand only matches a pure {@code
* IS NOT NULL} (or AND-of-{@code IS NOT NULL}) bucket-non-null filter; if {@code FilterMergeRule}
* runs first when a user {@code WHERE} sits below the synthetic {@code IS NOT NULL} filter that
* PPL emits as part of {@code dedup}, the two adjacent filters are merged into a single filter
* whose condition includes the user predicate, the simplify rule's predicate fails, no {@link
* org.opensearch.sql.calcite.plan.rel.LogicalDedup} is produced, and dedup pushdown to the
* OpenSearch storage engine is silently disabled. Using separate {@code addRuleInstance} calls
* (rather than {@code addRuleCollection}) enforces deterministic ordering: dedup simplification
* fires first against the original adjacent-filter shape, then any remaining adjacent filters are
* merged.
*/
private static final HepProgram HEP_PROGRAM =
new HepProgramBuilder().addRuleCollection(hepRuleList).build();
new HepProgramBuilder()
.addRuleInstance(PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE)
.addRuleInstance(FilterMergeRule.Config.DEFAULT.toRule())
.build();

public static RelNode optimize(RelNode plan, CalcitePlanContext context) {
Util.discard(context);
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -45,6 +45,7 @@
import org.opensearch.sql.calcite.CalcitePlanContext;
import org.opensearch.sql.calcite.CalciteRelNodeVisitor;
import org.opensearch.sql.calcite.SysLimit;
import org.opensearch.sql.calcite.utils.CalciteToolsHelper;
import org.opensearch.sql.common.setting.Settings;
import org.opensearch.sql.datasource.DataSourceService;
import org.opensearch.sql.exception.ExpressionEvaluationException;
Expand DownExpand Up@@ -110,6 +111,20 @@ public RelNode getRelNode(String ppl) {
return root;
}

/**
* Get the root RelNode of the given PPL query after running the production HEP program from
* {@code CalciteToolsHelper}. Use this in regression tests that exercise rules registered in the
* production HEP program (e.g. {@code PPLSimplifyDedupRule}) — those rules need to see the raw
* planner output, not the post-FilterMerge form returned by {@link #getRelNode(String)}.
*/
public RelNode getRelNodeAfterCalciteHep(String ppl) {
CalcitePlanContext context = createBuilderContext();
Query query = (Query) plan(pplParser, ppl);
planTransformer.analyze(query.getPlan(), context);
RelNode root = context.relBuilder.build();
return CalciteToolsHelper.optimize(root, context);
}

private RelNode mergeAdjacentFilters(RelNode relNode) {
HepProgram program =
new HepProgramBuilder().addRuleInstance(FilterMergeRule.Config.DEFAULT.toRule()).build();
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -5,9 +5,17 @@

package org.opensearch.sql.ppl.calcite;

import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;

import org.apache.calcite.plan.hep.HepPlanner;
import org.apache.calcite.plan.hep.HepProgram;
import org.apache.calcite.plan.hep.HepProgramBuilder;
import org.apache.calcite.rel.RelNode;
import org.apache.calcite.test.CalciteAssert;
import org.junit.Test;
import org.opensearch.sql.calcite.plan.rel.LogicalDedup;
import org.opensearch.sql.calcite.plan.rule.PPLSimplifyDedupRule;

public class CalcitePPLDedupTest extends CalcitePPLAbstractTest {

Expand DownExpand Up@@ -353,4 +361,72 @@ public void testSortFieldProjectedAwayBeforeDedup() {
+ " LogicalTableScan(table=[[scott, EMP]])\n";
verifyLogical(root, expectedLogical);
}

/**
* Regression test for issue #7: when a user {@code where} sits below {@code dedup}, the HEP
* program in {@code CalciteToolsHelper} must still produce a {@link LogicalDedup}. Before the
* fix, both rules were registered via {@code addRuleCollection}, so {@code FilterMergeRule} could
* fire ahead of {@code PPLSimplifyDedupRule} and merge the user predicate into the
* bucket-non-null filter; the simplify rule's bottom operand then rejected the merged condition
* (it only accepts pure {@code IS NOT NULL}/AND-of-{@code IS NOT NULL}), no {@code LogicalDedup}
* was produced, and dedup pushdown to the OpenSearch storage engine was silently disabled. The
* fix is to register the two rules with separate {@code addRuleInstance} calls in the order
* simplify-dedup first (to fixpoint), then filter-merge.
*/
@Test
public void testWhereThenDedupProducesLogicalDedup() {
// Use a where predicate on a DIFFERENT column from the dedup column. With the same column,
// Calcite's RexSimplify can fold AND(IS_NOT_NULL(x), >(x, c)) down to >(x, c), masking the
// bug. The issue's reproducer (where on @timestamp, dedup on namespace) hits this exact
// shape.
String ppl = "source=EMP | where SAL > 1000 | dedup 1 DEPTNO | fields DEPTNO";
RelNode optimized = getRelNodeAfterCalciteHep(ppl);
String optimizedPlan = optimized.explain();
assertTrue(
"where + dedup must produce a LogicalDedup so OpenSearch DedupPushdownRule can match;"
+ " actual plan was:\n"
+ optimizedPlan,
optimizedPlan.contains("LogicalDedup"));
// The window-form leftover would indicate the simplify rule did not fire — assert it is gone.
assertFalse(
"ROW_NUMBER window must be consumed by PPLSimplifyDedupRule when where + dedup are"
+ " combined; actual plan was:\n"
+ optimizedPlan,
optimizedPlan.contains("ROW_NUMBER"));
}

/**
* Adversarial regression test: simulates the pathological order described in issue #7 by forcing
* FilterMergeRule to run to fixpoint BEFORE PPLSimplifyDedupRule. This documents the failure mode
* the fix in {@code CalciteToolsHelper} prevents — once the bucket-non-null filter has been
* merged with the user {@code WHERE}, {@code mayBeFilterFromBucketNonNull} can never accept the
* combined condition, so {@code PPLSimplifyDedupRule} is permanently unable to produce a {@code
* LogicalDedup}. The production fix enforces order at the program level (sequential {@code
* addRuleInstance} calls), making this hazard unreachable.
*/
@Test
public void testFilterMergeBeforeSimplifyDedupBreaksPattern() {
String ppl = "source=EMP | where SAL > 1000 | dedup 1 DEPTNO | fields DEPTNO";
// getRelNode already runs FilterMergeRule on the raw plan, simulating the pathological
// schedule where FilterMergeRule fires before PPLSimplifyDedupRule.
RelNode mergedFirst = getRelNode(ppl);
HepProgram simplifyOnly =
new HepProgramBuilder().addRuleInstance(PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE).build();
HepPlanner planner = new HepPlanner(simplifyOnly);
planner.setRoot(mergedFirst);
RelNode result = planner.findBestExp();
String plan = result.explain();
assertFalse(
"If FilterMergeRule runs before PPLSimplifyDedupRule, the simplify rule must NOT recover"
+ " — the merged AND(IS_NOT_NULL, user_predicate) filter fails the bucket-non-null"
+ " predicate. This documents why the production HEP program enforces ordering via"
+ " separate addRuleInstance calls (PPLSimplifyDedupRule first, then FilterMergeRule)."
+ " Actual plan was:\n"
+ plan,
plan.contains("LogicalDedup"));
assertTrue(
"Plan should still contain ROW_NUMBER window form when simplify fails. Actual plan was:\n"
+ plan,
plan.contains("ROW_NUMBER"));
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line numberDiff line numberDiff line change
Expand Up@@ -36,7 +36,6 @@
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.time.Instant;
import java.util.List;
import java.util.Properties;
import java.util.function.Consumer;
import org.apache.calcite.adapter.enumerable.EnumerableConvention;
Expand All@@ -59,7 +58,6 @@
import org.apache.calcite.plan.Convention;
import org.apache.calcite.plan.RelOptCluster;
import org.apache.calcite.plan.RelOptPlanner;
import org.apache.calcite.plan.RelOptRule;
import org.apache.calcite.plan.RelOptSchema;
import org.apache.calcite.plan.RelOptTable;
import org.apache.calcite.plan.RelOptTable.ViewExpander;
Expand DownExpand Up@@ -551,12 +549,26 @@ public RelNode visit(TableScan scan) {
}
}

/** Try to optimize the plan by using HepPlanner */
private static final List<RelOptRule> hepRuleList =

Copy link
Copy Markdown
Owner

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ryan-gh-bot just for my knowledge why are we removing it?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

That hepRuleList field only existed as the argument to the old addRuleCollection(hepRuleList) call. Once the builder switched to two sequential addRuleInstance(...) calls (which take single rules directly), nothing else referenced the list, so keeping it would just be dead code — and dropping it also lets the java.util.List and org.apache.calcite.plan.RelOptRule imports go away. The Javadoc that used to sit on the field moved up onto HEP_PROGRAM, which is now the thing whose ordering actually matters.

List.of(FilterMergeRule.Config.DEFAULT.toRule(), PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE);

/**
* Try to optimize the plan by using HepPlanner.
*
* <p>Rule order matters: {@link PPLSimplifyDedupRule#DEDUP_SIMPLIFY_RULE} must run to fixpoint
* before {@link FilterMergeRule}. The simplify rule's bottom operand only matches a pure {@code
* IS NOT NULL} (or AND-of-{@code IS NOT NULL}) bucket-non-null filter; if {@code FilterMergeRule}
* runs first when a user {@code WHERE} sits below the synthetic {@code IS NOT NULL} filter that
* PPL emits as part of {@code dedup}, the two adjacent filters are merged into a single filter
* whose condition includes the user predicate, the simplify rule's predicate fails, no {@link
* org.opensearch.sql.calcite.plan.rel.LogicalDedup} is produced, and dedup pushdown to the
* OpenSearch storage engine is silently disabled. Using separate {@code addRuleInstance} calls
* (rather than {@code addRuleCollection}) enforces deterministic ordering: dedup simplification
* fires first against the original adjacent-filter shape, then any remaining adjacent filters are
* merged.
*/
private static final HepProgram HEP_PROGRAM =
new HepProgramBuilder().addRuleCollection(hepRuleList).build();
new HepProgramBuilder()
.addRuleInstance(PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE)
.addRuleInstance(FilterMergeRule.Config.DEFAULT.toRule())
.build();

public static RelNode optimize(RelNode plan, CalcitePlanContext context) {
Util.discard(context);
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -45,6 +45,7 @@
import org.opensearch.sql.calcite.CalcitePlanContext;
import org.opensearch.sql.calcite.CalciteRelNodeVisitor;
import org.opensearch.sql.calcite.SysLimit;
import org.opensearch.sql.calcite.utils.CalciteToolsHelper;
import org.opensearch.sql.common.setting.Settings;
import org.opensearch.sql.datasource.DataSourceService;
import org.opensearch.sql.exception.ExpressionEvaluationException;
Expand DownExpand Up@@ -110,6 +111,20 @@ public RelNode getRelNode(String ppl) {
return root;
}

/**
* Get the root RelNode of the given PPL query after running the production HEP program from
* {@code CalciteToolsHelper}. Use this in regression tests that exercise rules registered in the
* production HEP program (e.g. {@code PPLSimplifyDedupRule}) — those rules need to see the raw
* planner output, not the post-FilterMerge form returned by {@link #getRelNode(String)}.
*/
public RelNode getRelNodeAfterCalciteHep(String ppl) {
CalcitePlanContext context = createBuilderContext();
Query query = (Query) plan(pplParser, ppl);
planTransformer.analyze(query.getPlan(), context);
RelNode root = context.relBuilder.build();
return CalciteToolsHelper.optimize(root, context);
}

private RelNode mergeAdjacentFilters(RelNode relNode) {
HepProgram program =
new HepProgramBuilder().addRuleInstance(FilterMergeRule.Config.DEFAULT.toRule()).build();
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -5,9 +5,17 @@

package org.opensearch.sql.ppl.calcite;

import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;

import org.apache.calcite.plan.hep.HepPlanner;
import org.apache.calcite.plan.hep.HepProgram;
import org.apache.calcite.plan.hep.HepProgramBuilder;
import org.apache.calcite.rel.RelNode;
import org.apache.calcite.test.CalciteAssert;
import org.junit.Test;
import org.opensearch.sql.calcite.plan.rel.LogicalDedup;
import org.opensearch.sql.calcite.plan.rule.PPLSimplifyDedupRule;

public class CalcitePPLDedupTest extends CalcitePPLAbstractTest {

Expand DownExpand Up@@ -353,4 +361,72 @@ public void testSortFieldProjectedAwayBeforeDedup() {
+ " LogicalTableScan(table=[[scott, EMP]])\n";
verifyLogical(root, expectedLogical);
}

/**
* Regression test for issue #7: when a user {@code where} sits below {@code dedup}, the HEP
* program in {@code CalciteToolsHelper} must still produce a {@link LogicalDedup}. Before the
* fix, both rules were registered via {@code addRuleCollection}, so {@code FilterMergeRule} could
* fire ahead of {@code PPLSimplifyDedupRule} and merge the user predicate into the
* bucket-non-null filter; the simplify rule's bottom operand then rejected the merged condition
* (it only accepts pure {@code IS NOT NULL}/AND-of-{@code IS NOT NULL}), no {@code LogicalDedup}
* was produced, and dedup pushdown to the OpenSearch storage engine was silently disabled. The
* fix is to register the two rules with separate {@code addRuleInstance} calls in the order
* simplify-dedup first (to fixpoint), then filter-merge.
*/
@Test
public void testWhereThenDedupProducesLogicalDedup() {
// Use a where predicate on a DIFFERENT column from the dedup column. With the same column,
// Calcite's RexSimplify can fold AND(IS_NOT_NULL(x), >(x, c)) down to >(x, c), masking the
// bug. The issue's reproducer (where on @timestamp, dedup on namespace) hits this exact
// shape.
String ppl = "source=EMP | where SAL > 1000 | dedup 1 DEPTNO | fields DEPTNO";
RelNode optimized = getRelNodeAfterCalciteHep(ppl);
String optimizedPlan = optimized.explain();
assertTrue(
"where + dedup must produce a LogicalDedup so OpenSearch DedupPushdownRule can match;"
+ " actual plan was:\n"
+ optimizedPlan,
optimizedPlan.contains("LogicalDedup"));
// The window-form leftover would indicate the simplify rule did not fire — assert it is gone.
assertFalse(
"ROW_NUMBER window must be consumed by PPLSimplifyDedupRule when where + dedup are"
+ " combined; actual plan was:\n"
+ optimizedPlan,
optimizedPlan.contains("ROW_NUMBER"));
}

/**
* Adversarial regression test: simulates the pathological order described in issue #7 by forcing
* FilterMergeRule to run to fixpoint BEFORE PPLSimplifyDedupRule. This documents the failure mode
* the fix in {@code CalciteToolsHelper} prevents — once the bucket-non-null filter has been
* merged with the user {@code WHERE}, {@code mayBeFilterFromBucketNonNull} can never accept the
* combined condition, so {@code PPLSimplifyDedupRule} is permanently unable to produce a {@code
* LogicalDedup}. The production fix enforces order at the program level (sequential {@code
* addRuleInstance} calls), making this hazard unreachable.
*/
@Test
public void testFilterMergeBeforeSimplifyDedupBreaksPattern() {
String ppl = "source=EMP | where SAL > 1000 | dedup 1 DEPTNO | fields DEPTNO";
// getRelNode already runs FilterMergeRule on the raw plan, simulating the pathological
// schedule where FilterMergeRule fires before PPLSimplifyDedupRule.
RelNode mergedFirst = getRelNode(ppl);
HepProgram simplifyOnly =
new HepProgramBuilder().addRuleInstance(PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE).build();
HepPlanner planner = new HepPlanner(simplifyOnly);
planner.setRoot(mergedFirst);
RelNode result = planner.findBestExp();
String plan = result.explain();
assertFalse(
"If FilterMergeRule runs before PPLSimplifyDedupRule, the simplify rule must NOT recover"
+ " — the merged AND(IS_NOT_NULL, user_predicate) filter fails the bucket-non-null"
+ " predicate. This documents why the production HEP program enforces ordering via"
+ " separate addRuleInstance calls (PPLSimplifyDedupRule first, then FilterMergeRule)."
+ " Actual plan was:\n"
+ plan,
plan.contains("LogicalDedup"));
assertTrue(
"Plan should still contain ROW_NUMBER window form when simplify fails. Actual plan was:\n"
+ plan,
plan.contains("ROW_NUMBER"));
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line numberDiff line numberDiff line change
Expand Up@@ -36,7 +36,6 @@
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.time.Instant;
import java.util.List;
import java.util.Properties;
import java.util.function.Consumer;
import org.apache.calcite.adapter.enumerable.EnumerableConvention;
Expand All@@ -59,7 +58,6 @@
import org.apache.calcite.plan.Convention;
import org.apache.calcite.plan.RelOptCluster;
import org.apache.calcite.plan.RelOptPlanner;
import org.apache.calcite.plan.RelOptRule;
import org.apache.calcite.plan.RelOptSchema;
import org.apache.calcite.plan.RelOptTable;
import org.apache.calcite.plan.RelOptTable.ViewExpander;
Expand DownExpand Up@@ -551,12 +549,26 @@ public RelNode visit(TableScan scan) {
}
}

/** Try to optimize the plan by using HepPlanner */
private static final List<RelOptRule> hepRuleList =

Copy link
Copy Markdown
Owner

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ryan-gh-bot just for my knowledge why are we removing it?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

That hepRuleList field only existed as the argument to the old addRuleCollection(hepRuleList) call. Once the builder switched to two sequential addRuleInstance(...) calls (which take single rules directly), nothing else referenced the list, so keeping it would just be dead code — and dropping it also lets the java.util.List and org.apache.calcite.plan.RelOptRule imports go away. The Javadoc that used to sit on the field moved up onto HEP_PROGRAM, which is now the thing whose ordering actually matters.

List.of(FilterMergeRule.Config.DEFAULT.toRule(), PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE);

/**
* Try to optimize the plan by using HepPlanner.
*
* <p>Rule order matters: {@link PPLSimplifyDedupRule#DEDUP_SIMPLIFY_RULE} must run to fixpoint
* before {@link FilterMergeRule}. The simplify rule's bottom operand only matches a pure {@code
* IS NOT NULL} (or AND-of-{@code IS NOT NULL}) bucket-non-null filter; if {@code FilterMergeRule}
* runs first when a user {@code WHERE} sits below the synthetic {@code IS NOT NULL} filter that
* PPL emits as part of {@code dedup}, the two adjacent filters are merged into a single filter
* whose condition includes the user predicate, the simplify rule's predicate fails, no {@link
* org.opensearch.sql.calcite.plan.rel.LogicalDedup} is produced, and dedup pushdown to the
* OpenSearch storage engine is silently disabled. Using separate {@code addRuleInstance} calls
* (rather than {@code addRuleCollection}) enforces deterministic ordering: dedup simplification
* fires first against the original adjacent-filter shape, then any remaining adjacent filters are
* merged.
*/
private static final HepProgram HEP_PROGRAM =
new HepProgramBuilder().addRuleCollection(hepRuleList).build();
new HepProgramBuilder()
.addRuleInstance(PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE)
.addRuleInstance(FilterMergeRule.Config.DEFAULT.toRule())
.build();

public static RelNode optimize(RelNode plan, CalcitePlanContext context) {
Util.discard(context);
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -45,6 +45,7 @@
import org.opensearch.sql.calcite.CalcitePlanContext;
import org.opensearch.sql.calcite.CalciteRelNodeVisitor;
import org.opensearch.sql.calcite.SysLimit;
import org.opensearch.sql.calcite.utils.CalciteToolsHelper;
import org.opensearch.sql.common.setting.Settings;
import org.opensearch.sql.datasource.DataSourceService;
import org.opensearch.sql.exception.ExpressionEvaluationException;
Expand DownExpand Up@@ -110,6 +111,20 @@ public RelNode getRelNode(String ppl) {
return root;
}

/**
* Get the root RelNode of the given PPL query after running the production HEP program from
* {@code CalciteToolsHelper}. Use this in regression tests that exercise rules registered in the
* production HEP program (e.g. {@code PPLSimplifyDedupRule}) — those rules need to see the raw
* planner output, not the post-FilterMerge form returned by {@link #getRelNode(String)}.
*/
public RelNode getRelNodeAfterCalciteHep(String ppl) {
CalcitePlanContext context = createBuilderContext();
Query query = (Query) plan(pplParser, ppl);
planTransformer.analyze(query.getPlan(), context);
RelNode root = context.relBuilder.build();
return CalciteToolsHelper.optimize(root, context);
}

private RelNode mergeAdjacentFilters(RelNode relNode) {
HepProgram program =
new HepProgramBuilder().addRuleInstance(FilterMergeRule.Config.DEFAULT.toRule()).build();
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -5,9 +5,17 @@

package org.opensearch.sql.ppl.calcite;

import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;

import org.apache.calcite.plan.hep.HepPlanner;
import org.apache.calcite.plan.hep.HepProgram;
import org.apache.calcite.plan.hep.HepProgramBuilder;
import org.apache.calcite.rel.RelNode;
import org.apache.calcite.test.CalciteAssert;
import org.junit.Test;
import org.opensearch.sql.calcite.plan.rel.LogicalDedup;
import org.opensearch.sql.calcite.plan.rule.PPLSimplifyDedupRule;

public class CalcitePPLDedupTest extends CalcitePPLAbstractTest {

Expand DownExpand Up@@ -353,4 +361,72 @@ public void testSortFieldProjectedAwayBeforeDedup() {
+ " LogicalTableScan(table=[[scott, EMP]])\n";
verifyLogical(root, expectedLogical);
}

/**
* Regression test for issue #7: when a user {@code where} sits below {@code dedup}, the HEP
* program in {@code CalciteToolsHelper} must still produce a {@link LogicalDedup}. Before the
* fix, both rules were registered via {@code addRuleCollection}, so {@code FilterMergeRule} could
* fire ahead of {@code PPLSimplifyDedupRule} and merge the user predicate into the
* bucket-non-null filter; the simplify rule's bottom operand then rejected the merged condition
* (it only accepts pure {@code IS NOT NULL}/AND-of-{@code IS NOT NULL}), no {@code LogicalDedup}
* was produced, and dedup pushdown to the OpenSearch storage engine was silently disabled. The
* fix is to register the two rules with separate {@code addRuleInstance} calls in the order
* simplify-dedup first (to fixpoint), then filter-merge.
*/
@Test
public void testWhereThenDedupProducesLogicalDedup() {
// Use a where predicate on a DIFFERENT column from the dedup column. With the same column,
// Calcite's RexSimplify can fold AND(IS_NOT_NULL(x), >(x, c)) down to >(x, c), masking the
// bug. The issue's reproducer (where on @timestamp, dedup on namespace) hits this exact
// shape.
String ppl = "source=EMP | where SAL > 1000 | dedup 1 DEPTNO | fields DEPTNO";
RelNode optimized = getRelNodeAfterCalciteHep(ppl);
String optimizedPlan = optimized.explain();
assertTrue(
"where + dedup must produce a LogicalDedup so OpenSearch DedupPushdownRule can match;"
+ " actual plan was:\n"
+ optimizedPlan,
optimizedPlan.contains("LogicalDedup"));
// The window-form leftover would indicate the simplify rule did not fire — assert it is gone.
assertFalse(
"ROW_NUMBER window must be consumed by PPLSimplifyDedupRule when where + dedup are"
+ " combined; actual plan was:\n"
+ optimizedPlan,
optimizedPlan.contains("ROW_NUMBER"));
}

/**
* Adversarial regression test: simulates the pathological order described in issue #7 by forcing
* FilterMergeRule to run to fixpoint BEFORE PPLSimplifyDedupRule. This documents the failure mode
* the fix in {@code CalciteToolsHelper} prevents — once the bucket-non-null filter has been
* merged with the user {@code WHERE}, {@code mayBeFilterFromBucketNonNull} can never accept the
* combined condition, so {@code PPLSimplifyDedupRule} is permanently unable to produce a {@code
* LogicalDedup}. The production fix enforces order at the program level (sequential {@code
* addRuleInstance} calls), making this hazard unreachable.
*/
@Test
public void testFilterMergeBeforeSimplifyDedupBreaksPattern() {
String ppl = "source=EMP | where SAL > 1000 | dedup 1 DEPTNO | fields DEPTNO";
// getRelNode already runs FilterMergeRule on the raw plan, simulating the pathological
// schedule where FilterMergeRule fires before PPLSimplifyDedupRule.
RelNode mergedFirst = getRelNode(ppl);
HepProgram simplifyOnly =
new HepProgramBuilder().addRuleInstance(PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE).build();
HepPlanner planner = new HepPlanner(simplifyOnly);
planner.setRoot(mergedFirst);
RelNode result = planner.findBestExp();
String plan = result.explain();
assertFalse(
"If FilterMergeRule runs before PPLSimplifyDedupRule, the simplify rule must NOT recover"
+ " — the merged AND(IS_NOT_NULL, user_predicate) filter fails the bucket-non-null"
+ " predicate. This documents why the production HEP program enforces ordering via"
+ " separate addRuleInstance calls (PPLSimplifyDedupRule first, then FilterMergeRule)."
+ " Actual plan was:\n"
+ plan,
plan.contains("LogicalDedup"));
assertTrue(
"Plan should still contain ROW_NUMBER window form when simplify fails. Actual plan was:\n"
+ plan,
plan.contains("ROW_NUMBER"));
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line numberDiff line numberDiff line change
Expand Up@@ -36,7 +36,6 @@
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.time.Instant;
import java.util.List;
import java.util.Properties;
import java.util.function.Consumer;
import org.apache.calcite.adapter.enumerable.EnumerableConvention;
Expand All@@ -59,7 +58,6 @@
import org.apache.calcite.plan.Convention;
import org.apache.calcite.plan.RelOptCluster;
import org.apache.calcite.plan.RelOptPlanner;
import org.apache.calcite.plan.RelOptRule;
import org.apache.calcite.plan.RelOptSchema;
import org.apache.calcite.plan.RelOptTable;
import org.apache.calcite.plan.RelOptTable.ViewExpander;
Expand DownExpand Up@@ -551,12 +549,26 @@ public RelNode visit(TableScan scan) {
}
}

/** Try to optimize the plan by using HepPlanner */
private static final List<RelOptRule> hepRuleList =

Copy link
Copy Markdown
Owner

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ryan-gh-bot just for my knowledge why are we removing it?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

That hepRuleList field only existed as the argument to the old addRuleCollection(hepRuleList) call. Once the builder switched to two sequential addRuleInstance(...) calls (which take single rules directly), nothing else referenced the list, so keeping it would just be dead code — and dropping it also lets the java.util.List and org.apache.calcite.plan.RelOptRule imports go away. The Javadoc that used to sit on the field moved up onto HEP_PROGRAM, which is now the thing whose ordering actually matters.

List.of(FilterMergeRule.Config.DEFAULT.toRule(), PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE);

/**
* Try to optimize the plan by using HepPlanner.
*
* <p>Rule order matters: {@link PPLSimplifyDedupRule#DEDUP_SIMPLIFY_RULE} must run to fixpoint
* before {@link FilterMergeRule}. The simplify rule's bottom operand only matches a pure {@code
* IS NOT NULL} (or AND-of-{@code IS NOT NULL}) bucket-non-null filter; if {@code FilterMergeRule}
* runs first when a user {@code WHERE} sits below the synthetic {@code IS NOT NULL} filter that
* PPL emits as part of {@code dedup}, the two adjacent filters are merged into a single filter
* whose condition includes the user predicate, the simplify rule's predicate fails, no {@link
* org.opensearch.sql.calcite.plan.rel.LogicalDedup} is produced, and dedup pushdown to the
* OpenSearch storage engine is silently disabled. Using separate {@code addRuleInstance} calls
* (rather than {@code addRuleCollection}) enforces deterministic ordering: dedup simplification
* fires first against the original adjacent-filter shape, then any remaining adjacent filters are
* merged.
*/
private static final HepProgram HEP_PROGRAM =
new HepProgramBuilder().addRuleCollection(hepRuleList).build();
new HepProgramBuilder()
.addRuleInstance(PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE)
.addRuleInstance(FilterMergeRule.Config.DEFAULT.toRule())
.build();

public static RelNode optimize(RelNode plan, CalcitePlanContext context) {
Util.discard(context);
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -45,6 +45,7 @@
import org.opensearch.sql.calcite.CalcitePlanContext;
import org.opensearch.sql.calcite.CalciteRelNodeVisitor;
import org.opensearch.sql.calcite.SysLimit;
import org.opensearch.sql.calcite.utils.CalciteToolsHelper;
import org.opensearch.sql.common.setting.Settings;
import org.opensearch.sql.datasource.DataSourceService;
import org.opensearch.sql.exception.ExpressionEvaluationException;
Expand DownExpand Up@@ -110,6 +111,20 @@ public RelNode getRelNode(String ppl) {
return root;
}

/**
* Get the root RelNode of the given PPL query after running the production HEP program from
* {@code CalciteToolsHelper}. Use this in regression tests that exercise rules registered in the
* production HEP program (e.g. {@code PPLSimplifyDedupRule}) — those rules need to see the raw
* planner output, not the post-FilterMerge form returned by {@link #getRelNode(String)}.
*/
public RelNode getRelNodeAfterCalciteHep(String ppl) {
CalcitePlanContext context = createBuilderContext();
Query query = (Query) plan(pplParser, ppl);
planTransformer.analyze(query.getPlan(), context);
RelNode root = context.relBuilder.build();
return CalciteToolsHelper.optimize(root, context);
}

private RelNode mergeAdjacentFilters(RelNode relNode) {
HepProgram program =
new HepProgramBuilder().addRuleInstance(FilterMergeRule.Config.DEFAULT.toRule()).build();
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -5,9 +5,17 @@

package org.opensearch.sql.ppl.calcite;

import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;

import org.apache.calcite.plan.hep.HepPlanner;
import org.apache.calcite.plan.hep.HepProgram;
import org.apache.calcite.plan.hep.HepProgramBuilder;
import org.apache.calcite.rel.RelNode;
import org.apache.calcite.test.CalciteAssert;
import org.junit.Test;
import org.opensearch.sql.calcite.plan.rel.LogicalDedup;
import org.opensearch.sql.calcite.plan.rule.PPLSimplifyDedupRule;

public class CalcitePPLDedupTest extends CalcitePPLAbstractTest {

Expand DownExpand Up@@ -353,4 +361,72 @@ public void testSortFieldProjectedAwayBeforeDedup() {
+ " LogicalTableScan(table=[[scott, EMP]])\n";
verifyLogical(root, expectedLogical);
}

/**
* Regression test for issue #7: when a user {@code where} sits below {@code dedup}, the HEP
* program in {@code CalciteToolsHelper} must still produce a {@link LogicalDedup}. Before the
* fix, both rules were registered via {@code addRuleCollection}, so {@code FilterMergeRule} could
* fire ahead of {@code PPLSimplifyDedupRule} and merge the user predicate into the
* bucket-non-null filter; the simplify rule's bottom operand then rejected the merged condition
* (it only accepts pure {@code IS NOT NULL}/AND-of-{@code IS NOT NULL}), no {@code LogicalDedup}
* was produced, and dedup pushdown to the OpenSearch storage engine was silently disabled. The
* fix is to register the two rules with separate {@code addRuleInstance} calls in the order
* simplify-dedup first (to fixpoint), then filter-merge.
*/
@Test
public void testWhereThenDedupProducesLogicalDedup() {
// Use a where predicate on a DIFFERENT column from the dedup column. With the same column,
// Calcite's RexSimplify can fold AND(IS_NOT_NULL(x), >(x, c)) down to >(x, c), masking the
// bug. The issue's reproducer (where on @timestamp, dedup on namespace) hits this exact
// shape.
String ppl = "source=EMP | where SAL > 1000 | dedup 1 DEPTNO | fields DEPTNO";
RelNode optimized = getRelNodeAfterCalciteHep(ppl);
String optimizedPlan = optimized.explain();
assertTrue(
"where + dedup must produce a LogicalDedup so OpenSearch DedupPushdownRule can match;"
+ " actual plan was:\n"
+ optimizedPlan,
optimizedPlan.contains("LogicalDedup"));
// The window-form leftover would indicate the simplify rule did not fire — assert it is gone.
assertFalse(
"ROW_NUMBER window must be consumed by PPLSimplifyDedupRule when where + dedup are"
+ " combined; actual plan was:\n"
+ optimizedPlan,
optimizedPlan.contains("ROW_NUMBER"));
}

/**
* Adversarial regression test: simulates the pathological order described in issue #7 by forcing
* FilterMergeRule to run to fixpoint BEFORE PPLSimplifyDedupRule. This documents the failure mode
* the fix in {@code CalciteToolsHelper} prevents — once the bucket-non-null filter has been
* merged with the user {@code WHERE}, {@code mayBeFilterFromBucketNonNull} can never accept the
* combined condition, so {@code PPLSimplifyDedupRule} is permanently unable to produce a {@code
* LogicalDedup}. The production fix enforces order at the program level (sequential {@code
* addRuleInstance} calls), making this hazard unreachable.
*/
@Test
public void testFilterMergeBeforeSimplifyDedupBreaksPattern() {
String ppl = "source=EMP | where SAL > 1000 | dedup 1 DEPTNO | fields DEPTNO";
// getRelNode already runs FilterMergeRule on the raw plan, simulating the pathological
// schedule where FilterMergeRule fires before PPLSimplifyDedupRule.
RelNode mergedFirst = getRelNode(ppl);
HepProgram simplifyOnly =
new HepProgramBuilder().addRuleInstance(PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE).build();
HepPlanner planner = new HepPlanner(simplifyOnly);
planner.setRoot(mergedFirst);
RelNode result = planner.findBestExp();
String plan = result.explain();
assertFalse(
"If FilterMergeRule runs before PPLSimplifyDedupRule, the simplify rule must NOT recover"
+ " — the merged AND(IS_NOT_NULL, user_predicate) filter fails the bucket-non-null"
+ " predicate. This documents why the production HEP program enforces ordering via"
+ " separate addRuleInstance calls (PPLSimplifyDedupRule first, then FilterMergeRule)."
+ " Actual plan was:\n"
+ plan,
plan.contains("LogicalDedup"));
assertTrue(
"Plan should still contain ROW_NUMBER window form when simplify fails. Actual plan was:\n"
+ plan,
plan.contains("ROW_NUMBER"));
}
}
Loading