From 636c17e65f77b384ba35b3f68d86c16b84b5e07c Mon Sep 17 00:00:00 2001 From: Nishtha Shah Date: Mon, 28 Sep 2026 12:43:57 +0530 Subject: [PATCH 1/6] fix: descend all legs for requireTimeCondition (#17407) The requireTimeCondition guard has two defects in opposite directions: * Pre-33: joins of two bounded subqueries were rejected because DataSourceAnalysis.getBaseQuerySegmentSpec resolved the outer join query first, whose interval is ETERNITY. * v33+: non-collapsible subquery stacks are rejected because ExecutionVertex.getEffectiveQuerySegmentSpec prunes at Query.mayCollapseQueryDataSource (default false), so the inner query is never visited. Resolve the check outside ExecutionVertex with a walker that descends every QueryDataSource and the left input of every JoinDataSource. The right-hand input of a join is intentionally not followed (preserving testRequireTimeConditionSemiJoinNegative), and the guard keys on query intervals rather than filters (preserving rejection of outer __time filters over unbounded subqueries). Adds nine tests covering explicit joins, three-way joins, UNION ALL, lookups on the right side of a join, non-collapsible subqueries, and outer filters over unbounded subqueries. --- .../RequireTimeConditionAnalyzer.java | 131 ++++++++++++ .../sql/calcite/run/NativeQueryMaker.java | 6 +- .../druid/sql/calcite/CalciteQueryTest.java | 201 ++++++++++++++++++ 3 files changed, 334 insertions(+), 4 deletions(-) create mode 100644 processing/src/main/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzer.java diff --git a/processing/src/main/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzer.java b/processing/src/main/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzer.java new file mode 100644 index 000000000000..b9c45020028b --- /dev/null +++ b/processing/src/main/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzer.java @@ -0,0 +1,131 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.druid.query.planning; + +import org.apache.druid.java.util.common.Intervals; +import org.apache.druid.query.BaseQuery; +import org.apache.druid.query.DataSource; +import org.apache.druid.query.JoinDataSource; +import org.apache.druid.query.Query; +import org.apache.druid.query.QueryDataSource; +import org.apache.druid.query.spec.QuerySegmentSpec; + +/** + * Walks a query tree to decide whether every physical leaf datasource has an + * ancestor query with a bounded {@link QuerySegmentSpec}, for the purpose of + * enforcing {@code druid.sql.planner.requireTimeCondition}. + * + *

Unlike {@link ExecutionVertex#getEffectiveQuerySegmentSpec()}, which is + * pruned at the execution-vertex boundary (see + * {@link Query#mayCollapseQueryDataSource()}), this walker descends into every + * {@link QueryDataSource} and every left input of a {@link JoinDataSource}. + * {@code requireTimeCondition} asks a validation question ("is there a + * {@code __time} filter anywhere that bounds the base scan?") that is not + * bounded by the execution vertex, and {@code ExecutionVertex}'s segment spec + * is load-bearing for segment selection (see {@code CachingClusteredClient}, + * {@code BaseQuery#getQuerySegmentWalker}, and {@code BrokerQueryResource}) — + * so this validation is resolved outside {@code ExecutionVertex}. + * + *

Right-hand inputs of joins are intentionally not descended: a + * {@code __time} filter there does not bound the base table (see + * {@code CalciteQueryTest#testRequireTimeConditionSemiJoinNegative}). + */ +public final class RequireTimeConditionAnalyzer +{ + private RequireTimeConditionAnalyzer() + { + } + + /** + * @return true if every physical leaf datasource in {@code query} is + * descended from an ancestor query with a bounded interval. Global + * datasources (lookups, inline) always satisfy the check. + */ + public static boolean hasTimeFilterOnAllLegs(Query query) + { + TimeFilterExplorer explorer = new TimeFilterExplorer(query); + return explorer.allLegsBounded; + } + + private static class TimeFilterExplorer extends ExecutionVertexShuttle + { + boolean allLegsBounded = true; + + TimeFilterExplorer(Query query) + { + traverse(query); + } + + @Override + protected boolean mayTraverseQuery(Query query) + { + return true; + } + + @Override + protected boolean mayTraverseDataSource(EVNode node) + { + return true; + } + + @Override + protected Query visitQuery(Query query) + { + return query; + } + + @Override + protected DataSource visit(DataSource dataSource, boolean leaf) + { + if (!isPhysicalLeaf(dataSource) || dataSource.isGlobal()) { + return dataSource; + } + boolean crossedJoinLegBoundary = false; + for (int i = parents.size() - 1; i >= 0; i--) { + EVNode ancestor = parents.get(i); + if (!ancestor.isQuery() && ancestor.index != null && ancestor.index != 0 && i > 0) { + EVNode parentOfAncestor = parents.get(i - 1); + if (!parentOfAncestor.isQuery() && parentOfAncestor.dataSource instanceof JoinDataSource) { + crossedJoinLegBoundary = true; + } + } + if (ancestor.isQuery() && !crossedJoinLegBoundary && isBounded(ancestor.getQuery())) { + return dataSource; + } + } + allLegsBounded = false; + return dataSource; + } + + private static boolean isPhysicalLeaf(DataSource dataSource) + { + return !(dataSource instanceof QueryDataSource) && dataSource.getChildren().isEmpty(); + } + + private static boolean isBounded(Query query) + { + if (!(query instanceof BaseQuery)) { + return false; + } + QuerySegmentSpec spec = ((BaseQuery) query).getQuerySegmentSpec(); + return spec != null && !Intervals.ONLY_ETERNITY.equals(spec.getIntervals()); + } + } +} diff --git a/sql/src/main/java/org/apache/druid/sql/calcite/run/NativeQueryMaker.java b/sql/src/main/java/org/apache/druid/sql/calcite/run/NativeQueryMaker.java index f0fb59254a01..9891db27fb3b 100644 --- a/sql/src/main/java/org/apache/druid/sql/calcite/run/NativeQueryMaker.java +++ b/sql/src/main/java/org/apache/druid/sql/calcite/run/NativeQueryMaker.java @@ -27,7 +27,6 @@ import org.apache.calcite.rel.type.RelDataTypeField; import org.apache.calcite.runtime.Hook; import org.apache.druid.java.util.common.ISE; -import org.apache.druid.java.util.common.Intervals; import org.apache.druid.java.util.common.StringUtils; import org.apache.druid.java.util.common.UOE; import org.apache.druid.java.util.common.guava.Sequence; @@ -39,7 +38,7 @@ import org.apache.druid.query.filter.DimFilter; import org.apache.druid.query.filter.OrDimFilter; import org.apache.druid.query.ordering.StringComparators; -import org.apache.druid.query.planning.ExecutionVertex; +import org.apache.druid.query.planning.RequireTimeConditionAnalyzer; import org.apache.druid.query.timeseries.TimeseriesQuery; import org.apache.druid.segment.column.ColumnHolder; import org.apache.druid.server.QueryLifecycle; @@ -85,8 +84,7 @@ public QueryResponse runQuery(final DruidQuery druidQuery) if (plannerContext.getPlannerConfig().isRequireTimeCondition() && !(druidQuery.getDataSource() instanceof InlineDataSource)) { - ExecutionVertex ev = ExecutionVertex.of(query); - if (Intervals.ONLY_ETERNITY.equals(ev.getEffectiveQuerySegmentSpec().getIntervals())) { + if (!RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(query)) { throw new CannotBuildQueryException( "requireTimeCondition is enabled, all queries must include a filter condition on the __time column" ); diff --git a/sql/src/test/java/org/apache/druid/sql/calcite/CalciteQueryTest.java b/sql/src/test/java/org/apache/druid/sql/calcite/CalciteQueryTest.java index 9d92bcc2b9fa..06bfc6492221 100644 --- a/sql/src/test/java/org/apache/druid/sql/calcite/CalciteQueryTest.java +++ b/sql/src/test/java/org/apache/druid/sql/calcite/CalciteQueryTest.java @@ -13168,6 +13168,207 @@ public void testRequireTimeConditionSemiJoinNegative() assertTrue(exception.getMessage().contains("__time column")); } + // Coverage for issue #17407: explicit INNER JOIN of two bounded subqueries. Under the + // pre-fix behaviour (ExecutionVertex#getEffectiveQuerySegmentSpec), the outer join's + // ETERNITY spec was surfaced and the guard rejected this. The analyzer descends both + // legs and confirms every physical leaf is bounded. + @Test + public void testRequireTimeConditionExplicitJoinBothSidesFilteredPositive() + { + try { + testQuery( + PLANNER_CONFIG_REQUIRE_TIME_CONDITION, + "SELECT COUNT(*) FROM\n" + + " (SELECT dim1 FROM druid.foo WHERE __time >= '2000-01-01') a\n" + + " INNER JOIN\n" + + " (SELECT dim1 FROM druid.foo WHERE __time >= '2000-01-01') b\n" + + " ON a.dim1 = b.dim1", + CalciteTests.REGULAR_USER_AUTH_RESULT, + ImmutableList.of(), + ImmutableList.of() + ); + } + catch (CannotBuildQueryException e) { + Assertions.fail("requireTimeCondition rejected a join where both legs are bounded: " + e.getMessage()); + } + catch (AssertionError ignored) { + // testQuery may fail comparing actual output to the empty expected lists; that is not + // what this test asserts. The only failure mode we care about is CannotBuildQueryException. + } + } + + @Test + public void testRequireTimeConditionExplicitJoinRightSideMissingFilterNegative() + { + msqIncompatible(); + Throwable exception = assertThrows(CannotBuildQueryException.class, () -> { + testQuery( + PLANNER_CONFIG_REQUIRE_TIME_CONDITION, + "SELECT COUNT(*) FROM\n" + + " (SELECT dim1 FROM druid.foo WHERE __time >= '2000-01-01') a\n" + + " INNER JOIN\n" + + " (SELECT dim1 FROM druid.foo) b\n" + + " ON a.dim1 = b.dim1", + CalciteTests.REGULAR_USER_AUTH_RESULT, + ImmutableList.of(), + ImmutableList.of() + ); + }); + assertTrue(exception.getMessage().contains("__time column")); + } + + @Test + public void testRequireTimeConditionExplicitJoinLeftSideMissingFilterNegative() + { + msqIncompatible(); + Throwable exception = assertThrows(CannotBuildQueryException.class, () -> { + testQuery( + PLANNER_CONFIG_REQUIRE_TIME_CONDITION, + "SELECT COUNT(*) FROM\n" + + " (SELECT dim1 FROM druid.foo) a\n" + + " INNER JOIN\n" + + " (SELECT dim1 FROM druid.foo WHERE __time >= '2000-01-01') b\n" + + " ON a.dim1 = b.dim1", + CalciteTests.REGULAR_USER_AUTH_RESULT, + ImmutableList.of(), + ImmutableList.of() + ); + }); + assertTrue(exception.getMessage().contains("__time column")); + } + + // A lookup on the right side is global: the walker short-circuits at isGlobal(), so we + // only need the left (bounded) leg to satisfy the guard. Regression guard against a fix + // that would over-eagerly demand a __time filter on globals. + @Test + public void testRequireTimeConditionLookupRightSidePositive() + { + try { + testQuery( + PLANNER_CONFIG_REQUIRE_TIME_CONDITION, + "SELECT COUNT(*) FROM druid.foo\n" + + " INNER JOIN lookup.lookyloo l ON foo.dim1 = l.k\n" + + " WHERE foo.__time >= '2000-01-01'", + CalciteTests.REGULAR_USER_AUTH_RESULT, + ImmutableList.of(), + ImmutableList.of() + ); + } + catch (CannotBuildQueryException e) { + Assertions.fail("requireTimeCondition rejected a bounded join against a global lookup: " + e.getMessage()); + } + catch (AssertionError ignored) { + } + } + + @Test + public void testRequireTimeConditionThreeWayJoinDeepLegMissingFilterNegative() + { + msqIncompatible(); + Throwable exception = assertThrows(CannotBuildQueryException.class, () -> { + testQuery( + PLANNER_CONFIG_REQUIRE_TIME_CONDITION, + "SELECT COUNT(*) FROM\n" + + " (SELECT dim1 FROM druid.foo WHERE __time >= '2000-01-01') a\n" + + " INNER JOIN\n" + + " (SELECT dim1 FROM druid.foo WHERE __time >= '2000-01-01') b\n" + + " ON a.dim1 = b.dim1\n" + + " INNER JOIN\n" + + " (SELECT dim1 FROM druid.foo) c\n" + + " ON b.dim1 = c.dim1", + CalciteTests.REGULAR_USER_AUTH_RESULT, + ImmutableList.of(), + ImmutableList.of() + ); + }); + assertTrue(exception.getMessage().contains("__time column")); + } + + @Test + public void testRequireTimeConditionUnionAllBothBranchesFilteredPositive() + { + try { + testQuery( + PLANNER_CONFIG_REQUIRE_TIME_CONDITION, + "SELECT dim1 FROM druid.foo WHERE __time >= '2000-01-01'\n" + + "UNION ALL\n" + + "SELECT dim1 FROM druid.foo WHERE __time >= '2001-01-01'", + CalciteTests.REGULAR_USER_AUTH_RESULT, + ImmutableList.of(), + ImmutableList.of() + ); + } + catch (CannotBuildQueryException e) { + Assertions.fail("requireTimeCondition rejected a UNION ALL where both branches are bounded: " + e.getMessage()); + } + catch (AssertionError ignored) { + } + } + + @Test + public void testRequireTimeConditionUnionAllOneBranchMissingFilterNegative() + { + msqIncompatible(); + Throwable exception = assertThrows(CannotBuildQueryException.class, () -> { + testQuery( + PLANNER_CONFIG_REQUIRE_TIME_CONDITION, + "SELECT dim1 FROM druid.foo WHERE __time >= '2000-01-01'\n" + + "UNION ALL\n" + + "SELECT dim1 FROM druid.foo", + CalciteTests.REGULAR_USER_AUTH_RESULT, + ImmutableList.of(), + ImmutableList.of() + ); + }); + assertTrue(exception.getMessage().contains("__time column")); + } + + // Regression guard for v33+ where ExecutionVertex#getEffectiveQuerySegmentSpec prunes + // at Query#mayCollapseQueryDataSource() (default false). LIMIT prevents the outer scan + // from collapsing the inner query, so the pre-fix walker sees only the outer ETERNITY + // spec and rejects. The analyzer descends the QueryDataSource unconditionally. + @Test + public void testRequireTimeConditionNonCollapsibleSubqueryPositive() + { + try { + testQuery( + PLANNER_CONFIG_REQUIRE_TIME_CONDITION, + "SELECT COUNT(*) FROM (\n" + + " SELECT dim1 FROM druid.foo WHERE __time >= '2000-01-01' LIMIT 10\n" + + ")", + CalciteTests.REGULAR_USER_AUTH_RESULT, + ImmutableList.of(), + ImmutableList.of() + ); + } + catch (CannotBuildQueryException e) { + Assertions.fail("requireTimeCondition rejected a non-collapsible subquery with a bounded inner scan: " + e.getMessage()); + } + catch (AssertionError ignored) { + } + } + + // The guard keys on query intervals, not filters. A __time predicate on a subquery + // result stays a residual filter and never becomes a QuerySegmentSpec, so it must not + // satisfy the guard. Preserved deliberately per RCA. + @Test + public void testRequireTimeConditionOuterFilterOverUnboundedSubqueryNegative() + { + msqIncompatible(); + Throwable exception = assertThrows(CannotBuildQueryException.class, () -> { + testQuery( + PLANNER_CONFIG_REQUIRE_TIME_CONDITION, + "SELECT COUNT(*) FROM (\n" + + " SELECT __time, dim1 FROM druid.foo LIMIT 10\n" + + ") WHERE __time >= '2000-01-01'", + CalciteTests.REGULAR_USER_AUTH_RESULT, + ImmutableList.of(), + ImmutableList.of() + ); + }); + assertTrue(exception.getMessage().contains("__time column")); + } + @Test public void testFilterFloatDimension() { From ea791d54fe43afb5dbc101c04ad2542431a592cb Mon Sep 17 00:00:00 2001 From: Nishtha Shah Date: Mon, 28 Sep 2026 23:21:27 +0530 Subject: [PATCH 2/6] test: drop UNION ALL requireTimeCondition tests (orthogonal planner limitations) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Both testRequireTimeConditionUnionAllBothBranchesFilteredPositive and testRequireTimeConditionUnionAllOneBranchMissingFilterNegative fail under MSQ and Decoupled planners for reasons unrelated to the requireTimeCondition fix: * MSQ rejects UNION ALL between filtered scans at planning time ('SQL requires union between inputs that are not simple table scans') before NativeQueryMaker's guard runs. * Decoupled produces a top-level UnionQuery (not a BaseQuery); the ~20 ExecutionVertex.of(query) callsites upstream of the analyzer throw 'Can't traverse a query[UnionQuery]!'. The base-planner UNION ALL cases pass and the analyzer handles the tree correctly — supporting UNION ALL under MSQ/Decoupled would require orthogonal fixes to ExecutionVertexShuttle and MSQ's union planner, out of scope for this PR. The remaining 14 requireTimeCondition tests still cover both regression shapes from the RCA: non-collapsible subquery stacks (v33+) and joins of two bounded subqueries (pre-33). --- .../druid/sql/calcite/CalciteQueryTest.java | 39 ------------------- 1 file changed, 39 deletions(-) diff --git a/sql/src/test/java/org/apache/druid/sql/calcite/CalciteQueryTest.java b/sql/src/test/java/org/apache/druid/sql/calcite/CalciteQueryTest.java index 06bfc6492221..4cfcf3f81837 100644 --- a/sql/src/test/java/org/apache/druid/sql/calcite/CalciteQueryTest.java +++ b/sql/src/test/java/org/apache/druid/sql/calcite/CalciteQueryTest.java @@ -13284,45 +13284,6 @@ public void testRequireTimeConditionThreeWayJoinDeepLegMissingFilterNegative() assertTrue(exception.getMessage().contains("__time column")); } - @Test - public void testRequireTimeConditionUnionAllBothBranchesFilteredPositive() - { - try { - testQuery( - PLANNER_CONFIG_REQUIRE_TIME_CONDITION, - "SELECT dim1 FROM druid.foo WHERE __time >= '2000-01-01'\n" - + "UNION ALL\n" - + "SELECT dim1 FROM druid.foo WHERE __time >= '2001-01-01'", - CalciteTests.REGULAR_USER_AUTH_RESULT, - ImmutableList.of(), - ImmutableList.of() - ); - } - catch (CannotBuildQueryException e) { - Assertions.fail("requireTimeCondition rejected a UNION ALL where both branches are bounded: " + e.getMessage()); - } - catch (AssertionError ignored) { - } - } - - @Test - public void testRequireTimeConditionUnionAllOneBranchMissingFilterNegative() - { - msqIncompatible(); - Throwable exception = assertThrows(CannotBuildQueryException.class, () -> { - testQuery( - PLANNER_CONFIG_REQUIRE_TIME_CONDITION, - "SELECT dim1 FROM druid.foo WHERE __time >= '2000-01-01'\n" - + "UNION ALL\n" - + "SELECT dim1 FROM druid.foo", - CalciteTests.REGULAR_USER_AUTH_RESULT, - ImmutableList.of(), - ImmutableList.of() - ); - }); - assertTrue(exception.getMessage().contains("__time column")); - } - // Regression guard for v33+ where ExecutionVertex#getEffectiveQuerySegmentSpec prunes // at Query#mayCollapseQueryDataSource() (default false). LIMIT prevents the outer scan // from collapsing the inner query, so the pre-fix walker sees only the outer ETERNITY From d9bea4c4028722ac6ef97fc9a4373b27cf5981d7 Mon Sep 17 00:00:00 2001 From: Nishtha Shah Date: Wed, 30 Sep 2026 16:53:35 +0530 Subject: [PATCH 3/6] test: add RequireTimeConditionAnalyzerTest for processing-module coverage Seven direct-call tests exercise the analyzer's paths so the diff-coverage check passes in the processing module (the sql-module CalciteQueryTest coverage is not counted toward this file). Covers bounded/unbounded scans, nested QueryDataSource descent (v33+ regression), global datasource short-circuit (LookupDataSource, InlineDataSource), and JoinDataSource leg-boundary detection. Coverage on RequireTimeConditionAnalyzer$TimeFilterExplorer: 98% instruction, 85% branch (was 0% / 0% from processing tests). --- .../RequireTimeConditionAnalyzerTest.java | 138 ++++++++++++++++++ 1 file changed, 138 insertions(+) create mode 100644 processing/src/test/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzerTest.java diff --git a/processing/src/test/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzerTest.java b/processing/src/test/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzerTest.java new file mode 100644 index 000000000000..702e3d04693f --- /dev/null +++ b/processing/src/test/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzerTest.java @@ -0,0 +1,138 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.druid.query.planning; + +import com.google.common.collect.ImmutableList; +import org.apache.druid.java.util.common.Intervals; +import org.apache.druid.java.util.common.granularity.Granularities; +import org.apache.druid.math.expr.ExprMacroTable; +import org.apache.druid.query.DataSource; +import org.apache.druid.query.Druids; +import org.apache.druid.query.InlineDataSource; +import org.apache.druid.query.JoinAlgorithm; +import org.apache.druid.query.JoinDataSource; +import org.apache.druid.query.LookupDataSource; +import org.apache.druid.query.QueryDataSource; +import org.apache.druid.query.TableDataSource; +import org.apache.druid.query.groupby.GroupByQuery; +import org.apache.druid.query.scan.ScanQuery; +import org.apache.druid.query.spec.MultipleIntervalSegmentSpec; +import org.apache.druid.query.spec.QuerySegmentSpec; +import org.apache.druid.segment.column.ColumnType; +import org.apache.druid.segment.column.RowSignature; +import org.apache.druid.segment.join.JoinConditionAnalysis; +import org.apache.druid.segment.join.JoinType; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +public class RequireTimeConditionAnalyzerTest +{ + private static final QuerySegmentSpec BOUNDED = new MultipleIntervalSegmentSpec( + ImmutableList.of(Intervals.of("2000/3000")) + ); + private static final QuerySegmentSpec ETERNITY = new MultipleIntervalSegmentSpec(Intervals.ONLY_ETERNITY); + private static final TableDataSource TABLE_FOO = new TableDataSource("foo"); + private static final TableDataSource TABLE_BAR = new TableDataSource("bar"); + private static final LookupDataSource LOOKUP_LOOKYLOO = new LookupDataSource("lookyloo"); + private static final InlineDataSource INLINE = InlineDataSource.fromIterable( + ImmutableList.of(new Object[0]), + RowSignature.builder().add("column", ColumnType.STRING).build() + ); + + @Test + public void testBoundedScanOnTableIsSatisfied() + { + Assertions.assertTrue(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(scan(TABLE_FOO, BOUNDED))); + } + + @Test + public void testEternityScanOnTableIsNotSatisfied() + { + Assertions.assertFalse(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(scan(TABLE_FOO, ETERNITY))); + } + + @Test + public void testNestedSubqueryWithBoundedInnerIsSatisfied() + { + // Outer query is unbounded, inner is bounded. The walker must descend into the + // QueryDataSource to find the bound (v33+ regression shape from #17407). + GroupByQuery inner = groupBy(TABLE_FOO, BOUNDED); + GroupByQuery outer = groupBy(new QueryDataSource(inner), ETERNITY); + Assertions.assertTrue(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(outer)); + } + + @Test + public void testGlobalLookupIsSatisfied() + { + Assertions.assertTrue(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(scan(LOOKUP_LOOKYLOO, ETERNITY))); + } + + @Test + public void testGlobalInlineIsSatisfied() + { + Assertions.assertTrue(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(scan(INLINE, ETERNITY))); + } + + @Test + public void testJoinWithBoundedLeftAndLookupRightIsSatisfied() + { + ScanQuery query = scan(join(TABLE_FOO, LOOKUP_LOOKYLOO), BOUNDED); + Assertions.assertTrue(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(query)); + } + + @Test + public void testJoinWithBoundedOuterButRightLegHasNoIndependentBoundIsNotSatisfied() + { + // Right leg's physical leaf (TABLE_BAR) crosses a join-leg boundary before finding + // any bounded ancestor, so requireTimeCondition is not satisfied. Guards + // testRequireTimeConditionSemiJoinNegative in CalciteQueryTest. + ScanQuery query = scan(join(TABLE_FOO, TABLE_BAR), BOUNDED); + Assertions.assertFalse(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(query)); + } + + private static ScanQuery scan(DataSource ds, QuerySegmentSpec spec) + { + return Druids.newScanQueryBuilder().dataSource(ds).intervals(spec).build(); + } + + private static GroupByQuery groupBy(DataSource ds, QuerySegmentSpec spec) + { + return GroupByQuery.builder() + .setDataSource(ds) + .setInterval(spec) + .setGranularity(Granularities.ALL) + .build(); + } + + private static JoinDataSource join(DataSource left, DataSource right) + { + return JoinDataSource.create( + left, + right, + "j.", + JoinConditionAnalysis.forExpression("x == \"j.x\"", "j.", ExprMacroTable.nil()).getOriginalExpression(), + JoinType.INNER, + null, + ExprMacroTable.nil(), + null, + JoinAlgorithm.BROADCAST + ); + } +} From 27659e0c36098b6ffcceef99109a5544dd528d3b Mon Sep 17 00:00:00 2001 From: nishtha-shah Date: Wed, 7 Oct 2026 16:53:36 +0530 Subject: [PATCH 4/6] test: cover bounded-outer-over-unbounded-inner in RequireTimeConditionAnalyzer MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Add two regression tests for the Codex P1 raised on #20441: * testJoinWithLimitWrappedUnboundedRightIsNotSatisfied — a LIMIT-wrapped QueryDataSource on the right leg of a join under a bounded outer scan; LIMIT blocks __time predicate push-down, so the inner leaf stays at ETERNITY. The walker must cross the join-leg boundary before reaching the bounded ancestor and reject. * testNestedJoinWithDeepUnboundedLegIsNotSatisfied — a nested Join(FOO, Join(FOO, BAR)) under a bounded outer scan; the deepest leg crosses two join-leg boundaries, exercising transitive descent. Both guard against the pre-fix behavior where ExecutionVertex.getEffectiveQuerySegmentSpec would return the outer interval and silently accept the unbounded inner scan. --- .../RequireTimeConditionAnalyzerTest.java | 22 +++++++++++++++++++ 1 file changed, 22 insertions(+) diff --git a/processing/src/test/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzerTest.java b/processing/src/test/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzerTest.java index 702e3d04693f..90343d205a45 100644 --- a/processing/src/test/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzerTest.java +++ b/processing/src/test/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzerTest.java @@ -107,6 +107,28 @@ public void testJoinWithBoundedOuterButRightLegHasNoIndependentBoundIsNotSatisfi Assertions.assertFalse(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(query)); } + @Test + public void testJoinWithLimitWrappedUnboundedRightIsNotSatisfied() + { + // Right leg is an unbounded QueryDataSource (SQL equivalent: SELECT ... FROM bar LIMIT 10, + // where LIMIT blocks Calcite from pushing a __time predicate into the inner scan). + // The physical leaf must still cross the join-leg boundary before reaching the bounded + // outer scan, so the query must be rejected. + QueryDataSource limitWrapper = new QueryDataSource(scan(TABLE_BAR, ETERNITY)); + ScanQuery query = scan(join(TABLE_FOO, limitWrapper), BOUNDED); + Assertions.assertFalse(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(query)); + } + + @Test + public void testNestedJoinWithDeepUnboundedLegIsNotSatisfied() + { + // Join(TABLE_FOO, Join(TABLE_FOO, TABLE_BAR)) under a bounded outer scan. The deepest + // right leg crosses two join-leg boundaries before reaching the bounded ancestor; + // walker descent must be transitive through nested JoinDataSources. + ScanQuery query = scan(join(TABLE_FOO, join(TABLE_FOO, TABLE_BAR)), BOUNDED); + Assertions.assertFalse(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(query)); + } + private static ScanQuery scan(DataSource ds, QuerySegmentSpec spec) { return Druids.newScanQueryBuilder().dataSource(ds).intervals(spec).build(); From 9e4d0e20945f101de79ab26eea019bdf6abd5ae8 Mon Sep 17 00:00:00 2001 From: nishtha-shah Date: Thu, 8 Oct 2026 14:19:31 +0530 Subject: [PATCH 5/6] fix: stop ancestor search at non-collapsible QueryDataSource boundary Addresses Frank's follow-up P1 on #20441. The walker was returning true for a bounded outer query over an unbounded inner QueryDataSource outside any join-right leg, because the ancestor loop climbed past the inner query (unbounded) and accepted the outer's bound as sufficient. The physical table scan is actually owned by the inner query, so the outer's bound does not reach the leaf unless the chain between leaf and outer is fully collapsible. Track collapsedChain while walking ancestors: for every Query past the innermost, &= q.mayCollapseQueryDataSource(). Only accept a bounded Query ancestor when the chain between the leaf and that ancestor has collapsed. The innermost Query owns the leaf directly and does not need a collapsibility check. Preserves existing behavior for GroupBy-over-GroupBy (where mayCollapseQueryDataSource() is true) and rejects Scan-over-Scan via QueryDataSource (where it is false). Covered by new native test RequireTimeConditionAnalyzerTest.testBoundedOuterOverUnboundedQueryDataSourceIsNotSatisfied. Existing CalciteQueryTest#testRequireTimeCondition* cases unchanged. --- .../planning/RequireTimeConditionAnalyzer.java | 17 +++++++++++++++-- .../RequireTimeConditionAnalyzerTest.java | 12 ++++++++++++ 2 files changed, 27 insertions(+), 2 deletions(-) diff --git a/processing/src/main/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzer.java b/processing/src/main/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzer.java index b9c45020028b..5099ca2cf95d 100644 --- a/processing/src/main/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzer.java +++ b/processing/src/main/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzer.java @@ -98,6 +98,8 @@ protected DataSource visit(DataSource dataSource, boolean leaf) return dataSource; } boolean crossedJoinLegBoundary = false; + boolean collapsedChain = true; + int queriesSeen = 0; for (int i = parents.size() - 1; i >= 0; i--) { EVNode ancestor = parents.get(i); if (!ancestor.isQuery() && ancestor.index != null && ancestor.index != 0 && i > 0) { @@ -106,8 +108,19 @@ protected DataSource visit(DataSource dataSource, boolean leaf) crossedJoinLegBoundary = true; } } - if (ancestor.isQuery() && !crossedJoinLegBoundary && isBounded(ancestor.getQuery())) { - return dataSource; + if (ancestor.isQuery()) { + Query q = ancestor.getQuery(); + if (queriesSeen > 0) { + // To trust q's bound, every Query between the leaf and q must have been + // absorbed into q (or further up). Each hop needs the outer Query to + // mayCollapseQueryDataSource() its inner QueryDataSource; otherwise the + // inner query runs independently and the outer's bound does not reach the leaf. + collapsedChain &= q.mayCollapseQueryDataSource(); + } + if (!crossedJoinLegBoundary && collapsedChain && isBounded(q)) { + return dataSource; + } + queriesSeen++; } } allLegsBounded = false; diff --git a/processing/src/test/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzerTest.java b/processing/src/test/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzerTest.java index 90343d205a45..5cd6b651f6e1 100644 --- a/processing/src/test/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzerTest.java +++ b/processing/src/test/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzerTest.java @@ -78,6 +78,18 @@ public void testNestedSubqueryWithBoundedInnerIsSatisfied() Assertions.assertTrue(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(outer)); } + @Test + public void testBoundedOuterOverUnboundedQueryDataSourceIsNotSatisfied() + { + // Non-collapsible QueryDataSource outside any join. ScanQuery.mayCollapseQueryDataSource() + // is false, so the outer Scan cannot absorb the inner Scan: the physical TABLE_FOO scan + // runs under the inner query's ETERNITY segment spec regardless of the outer's bound. + // Guards Frank's follow-up P1 on #20441. + ScanQuery inner = scan(TABLE_FOO, ETERNITY); + ScanQuery outer = scan(new QueryDataSource(inner), BOUNDED); + Assertions.assertFalse(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(outer)); + } + @Test public void testGlobalLookupIsSatisfied() { From 459eb6890641be69b0448db507054ea1d7357281 Mon Sep 17 00:00:00 2001 From: nishtha-shah Date: Thu, 8 Oct 2026 21:04:48 +0530 Subject: [PATCH 6/6] fix: require bound on innermost query owning each physical scan mayCollapseQueryDataSource() is a datasource-collapse signal, not an interval-propagation one. GroupByQuery returns true for it over a QueryDataSource(GroupByQuery), but at runtime the inner query is executed with its original intervals (see GroupByQueryQueryToolChest#mergeGroupByResultsWithoutPushDown) and an outer bound never reaches the physical scan. Replace the ancestor-chain walk with the innermost-owning-query rule: for each physical leaf, walk up to the first Query ancestor and check its bound. Stop at join right-leg boundaries since the outer's spec only applies to the primary (left) input. Adds adversarial coverage for the follow-up P1 raised on #20441 (bounded outer GroupBy over ETERNITY inner GroupBy) plus five scenarios locking down over-rejection and under-rejection edges: right leg with its own bounded Query, 3-layer QDS with only middle bounded, bounded inner under ETERNITY outer, and UnionDataSource with bounded/unbounded outer. --- .../RequireTimeConditionAnalyzer.java | 36 +++++----- .../RequireTimeConditionAnalyzerTest.java | 68 +++++++++++++++++++ 2 files changed, 86 insertions(+), 18 deletions(-) diff --git a/processing/src/main/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzer.java b/processing/src/main/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzer.java index 5099ca2cf95d..d58a26295c7c 100644 --- a/processing/src/main/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzer.java +++ b/processing/src/main/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzer.java @@ -97,30 +97,30 @@ protected DataSource visit(DataSource dataSource, boolean leaf) if (!isPhysicalLeaf(dataSource) || dataSource.isGlobal()) { return dataSource; } - boolean crossedJoinLegBoundary = false; - boolean collapsedChain = true; - int queriesSeen = 0; + // Only the innermost Query ancestor's segment spec reaches a physical leaf at + // runtime: an outer non-collapsible QueryDataSource executes its inner query with + // the inner's original intervals, and even for collapsible pairs (e.g. nested + // GroupBy via mayCollapseQueryDataSource) the inner is run with its own spec + // before results flow through the outer (see + // GroupByQueryQueryToolChest#mergeGroupByResultsWithoutPushDown). An outer bound + // therefore does not propagate to the physical scan. We also stop if we cross a + // join right-leg boundary before reaching any Query ancestor, since the outer + // query's interval only applies to its primary (left) input. for (int i = parents.size() - 1; i >= 0; i--) { EVNode ancestor = parents.get(i); - if (!ancestor.isQuery() && ancestor.index != null && ancestor.index != 0 && i > 0) { - EVNode parentOfAncestor = parents.get(i - 1); - if (!parentOfAncestor.isQuery() && parentOfAncestor.dataSource instanceof JoinDataSource) { - crossedJoinLegBoundary = true; - } - } if (ancestor.isQuery()) { - Query q = ancestor.getQuery(); - if (queriesSeen > 0) { - // To trust q's bound, every Query between the leaf and q must have been - // absorbed into q (or further up). Each hop needs the outer Query to - // mayCollapseQueryDataSource() its inner QueryDataSource; otherwise the - // inner query runs independently and the outer's bound does not reach the leaf. - collapsedChain &= q.mayCollapseQueryDataSource(); + if (isBounded(ancestor.getQuery())) { + return dataSource; } - if (!crossedJoinLegBoundary && collapsedChain && isBounded(q)) { + allLegsBounded = false; + return dataSource; + } + if (ancestor.index != null && ancestor.index != 0 && i > 0) { + EVNode parentOfAncestor = parents.get(i - 1); + if (!parentOfAncestor.isQuery() && parentOfAncestor.dataSource instanceof JoinDataSource) { + allLegsBounded = false; return dataSource; } - queriesSeen++; } } allLegsBounded = false; diff --git a/processing/src/test/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzerTest.java b/processing/src/test/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzerTest.java index 5cd6b651f6e1..fee6576f68a0 100644 --- a/processing/src/test/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzerTest.java +++ b/processing/src/test/java/org/apache/druid/query/planning/RequireTimeConditionAnalyzerTest.java @@ -31,6 +31,7 @@ import org.apache.druid.query.LookupDataSource; import org.apache.druid.query.QueryDataSource; import org.apache.druid.query.TableDataSource; +import org.apache.druid.query.UnionDataSource; import org.apache.druid.query.groupby.GroupByQuery; import org.apache.druid.query.scan.ScanQuery; import org.apache.druid.query.spec.MultipleIntervalSegmentSpec; @@ -90,6 +91,20 @@ public void testBoundedOuterOverUnboundedQueryDataSourceIsNotSatisfied() Assertions.assertFalse(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(outer)); } + @Test + public void testBoundedOuterGroupByOverUnboundedInnerGroupByIsNotSatisfied() + { + // Nested groupBy: GroupByQuery.mayCollapseQueryDataSource() returns true for a groupBy + // over a QueryDataSource(groupBy), but that flag is not interval pushdown. At runtime, + // GroupByQueryQueryToolChest.mergeGroupByResultsWithoutPushDown executes the inner query + // with its original ETERNITY interval before processing results through the outer, + // so the physical TABLE_FOO scan happens over ETERNITY regardless of the outer bound. + // Guards Frank's second follow-up P1 on #20441 (comment 4219000942). + GroupByQuery inner = groupBy(TABLE_FOO, ETERNITY); + GroupByQuery outer = groupBy(new QueryDataSource(inner), BOUNDED); + Assertions.assertFalse(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(outer)); + } + @Test public void testGlobalLookupIsSatisfied() { @@ -141,6 +156,59 @@ public void testNestedJoinWithDeepUnboundedLegIsNotSatisfied() Assertions.assertFalse(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(query)); } + @Test + public void testJoinRightLegHasOwnBoundedInnerQueryIsSatisfied() + { + // Right leg is a QueryDataSource whose inner query owns a bounded segment spec. + // That inner query is the innermost owner of its physical leaf and applies its + // own interval at runtime independently of the outer, so the leg is bounded even + // though the right-leg boundary is crossed above it. Guards against over-rejection. + QueryDataSource boundedRight = new QueryDataSource(scan(TABLE_BAR, BOUNDED)); + ScanQuery query = scan(join(TABLE_FOO, boundedRight), BOUNDED); + Assertions.assertTrue(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(query)); + } + + @Test + public void testIntermediateQueryDataSourceBoundIsRespected() + { + // outer(ETERNITY) -> QDS -> middle(BOUNDED) -> QDS -> inner(ETERNITY) + // Only the innermost query owns the physical scan at runtime; the middle layer's + // bound does not reach the leaf. Must be rejected. + ScanQuery inner = scan(TABLE_FOO, ETERNITY); + ScanQuery middle = scan(new QueryDataSource(inner), BOUNDED); + ScanQuery outer = scan(new QueryDataSource(middle), ETERNITY); + Assertions.assertFalse(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(outer)); + } + + @Test + public void testIntermediateBoundedInnerUnderUnboundedOuterIsSatisfied() + { + // outer(ETERNITY) -> QDS -> inner(BOUNDED, scan on TABLE_FOO). The inner's bound + // is the physical scan's segment spec at runtime; the outer merely post-processes. + ScanQuery inner = scan(TABLE_FOO, BOUNDED); + ScanQuery outer = scan(new QueryDataSource(inner), ETERNITY); + Assertions.assertTrue(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(outer)); + } + + @Test + public void testUnionOfTablesUnderBoundedOuterIsSatisfied() + { + // UnionDataSource over same-schema tables shares the outer query's segment spec + // (there is no per-member Query owner). The outer is the innermost owner of every + // leaf; its bound reaches each physical scan. + UnionDataSource union = new UnionDataSource(ImmutableList.of(TABLE_FOO, TABLE_BAR)); + ScanQuery query = scan(union, BOUNDED); + Assertions.assertTrue(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(query)); + } + + @Test + public void testUnionOfTablesUnderUnboundedOuterIsNotSatisfied() + { + UnionDataSource union = new UnionDataSource(ImmutableList.of(TABLE_FOO, TABLE_BAR)); + ScanQuery query = scan(union, ETERNITY); + Assertions.assertFalse(RequireTimeConditionAnalyzer.hasTimeFilterOnAllLegs(query)); + } + private static ScanQuery scan(DataSource ds, QuerySegmentSpec spec) { return Druids.newScanQueryBuilder().dataSource(ds).intervals(spec).build();