clintropolis commented on code in PR #19830: URL: https://github.com/apache/druid/pull/19830#discussion_r4029978484
########## sql/src/test/java/org/apache/druid/sql/calcite/CalciteCatalogDdlTest.java: ########## @@ -0,0 +1,1089 @@ +/* + * 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.sql.calcite; + +import com.google.common.collect.ImmutableList; +import com.google.common.collect.ImmutableMap; +import com.google.common.collect.ImmutableSet; +import org.apache.calcite.sql.SqlExplain; +import org.apache.calcite.sql.SqlExplainFormat; +import org.apache.calcite.sql.SqlExplainLevel; +import org.apache.calcite.sql.SqlNode; +import org.apache.calcite.sql.parser.SqlParserPos; +import org.apache.druid.catalog.model.ClusteredValueGroupsBaseTableMetadata; +import org.apache.druid.catalog.model.ColumnSpec; +import org.apache.druid.catalog.model.DatasourceBaseTableMetadata; +import org.apache.druid.catalog.model.DatasourceProjectionMetadata; +import org.apache.druid.catalog.model.TableId; +import org.apache.druid.catalog.model.TableMetadata; +import org.apache.druid.catalog.model.TableSpec; +import org.apache.druid.catalog.model.table.ClusterKeySpec; +import org.apache.druid.catalog.model.table.DatasourceDefn; +import org.apache.druid.data.input.impl.AggregateProjectionSpec; +import org.apache.druid.error.DruidException; +import org.apache.druid.java.util.common.granularity.Granularities; +import org.apache.druid.server.security.Action; +import org.apache.druid.server.security.AuthConfig; +import org.apache.druid.server.security.AuthenticationResult; +import org.apache.druid.server.security.ForbiddenException; +import org.apache.druid.server.security.Resource; +import org.apache.druid.server.security.ResourceAction; +import org.apache.druid.server.security.ResourceType; +import org.apache.druid.sql.DirectStatement; +import org.apache.druid.sql.SqlQueryPlus; +import org.apache.druid.sql.calcite.CalciteCatalogDdlTest.CatalogDdlComponentSupplier; +import org.apache.druid.sql.calcite.parser.DruidSqlParser; +import org.apache.druid.sql.calcite.planner.CatalogTableWriter; +import org.apache.druid.sql.calcite.planner.DruidPlanner; +import org.apache.druid.sql.calcite.planner.PlannerConfig; +import org.apache.druid.sql.calcite.planner.PlannerFactory; +import org.apache.druid.sql.calcite.util.CalciteTests; +import org.apache.druid.sql.calcite.util.SqlTestFramework.StandardComponentSupplier; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import javax.annotation.Nullable; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Tests that catalog DDL statements plan into the catalog operations they claim to, using a writer that records + * calls instead of contacting a Coordinator. + */ [email protected](CatalogDdlComponentSupplier.class) +public class CalciteCatalogDdlTest extends BaseCalciteQueryTest +{ + private static final RecordingCatalogTableWriter WRITER = new RecordingCatalogTableWriter(); + + public static class CatalogDdlComponentSupplier extends StandardComponentSupplier + { + public CatalogDdlComponentSupplier(TempDirProducer tempFolderProducer) + { + super(tempFolderProducer); + } + + @Override + public CatalogTableWriter createCatalogTableWriter() + { + return WRITER; + } + } + + @BeforeEach + public void resetWriter() + { + WRITER.reset(); + } + + @Test + public void testCreateTable() + { + execute("CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT)"); + + assertEquals(1, WRITER.calls.size()); + final RecordingCatalogTableWriter.Call call = WRITER.calls.get(0); + assertEquals("createTable", call.operation); + assertEquals(TableId.datasource("tbl"), call.tableId); + assertEquals(DatasourceDefn.TABLE_TYPE, call.spec.type()); + assertEquals(ImmutableMap.of(), call.spec.properties()); + assertEquals( + ImmutableList.of( + new ColumnSpec("__time", "TIMESTAMP", null), + new ColumnSpec("page", "VARCHAR", null), + new ColumnSpec("cnt", "BIGINT", null) + ), + call.spec.columns() + ); + assertFalse(call.ifNotExists); + assertFalse(call.replace); + } + + @Test + public void testCreateTableWithPartitioningAndClustering() + { + execute("CREATE TABLE tbl (page VARCHAR, cnt BIGINT) PARTITIONED BY DAY CLUSTERED BY page, cnt"); + + final TableSpec spec = WRITER.calls.get(0).spec; + assertEquals("P1D", spec.properties().get(DatasourceDefn.SEGMENT_GRANULARITY_PROPERTY)); + assertEquals( + ImmutableList.of(new ClusterKeySpec("page", false), new ClusterKeySpec("cnt", false)), + spec.properties().get(DatasourceDefn.CLUSTER_KEYS_PROPERTY) + ); + } + + @Test + public void testCreateTablePartitionedByAll() + { + execute("CREATE TABLE tbl (page VARCHAR) PARTITIONED BY ALL TIME"); + assertEquals("ALL", WRITER.calls.get(0).spec.properties().get(DatasourceDefn.SEGMENT_GRANULARITY_PROPERTY)); + } + + @Test + public void testCreateTableTypeCanonicalization() + { + execute( + "CREATE TABLE tbl (a CHAR, b INTEGER, c REAL, d DOUBLE, e VARCHAR ARRAY, f TYPE('complex<json>'))" + ); + assertEquals( + ImmutableList.of( + new ColumnSpec("a", "VARCHAR", null), + new ColumnSpec("b", "BIGINT", null), + new ColumnSpec("c", "FLOAT", null), + new ColumnSpec("d", "DOUBLE", null), + new ColumnSpec("e", "VARCHAR ARRAY", null), + new ColumnSpec("f", "COMPLEX<json>", null) + ), + WRITER.calls.get(0).spec.columns() + ); + } + + @Test + public void testCreateTableFlags() + { + execute("CREATE OR REPLACE TABLE tbl (a VARCHAR)"); + assertTrue(WRITER.calls.get(0).replace); + + WRITER.reset(); + execute("CREATE TABLE IF NOT EXISTS tbl (a VARCHAR)"); + assertTrue(WRITER.calls.get(0).ifNotExists); + } + + @Test + public void testCreateTableInDruidSchema() + { + execute("CREATE TABLE druid.tbl (a VARCHAR)"); + assertEquals(TableId.datasource("tbl"), WRITER.calls.get(0).tableId); + } + + @Test + public void testResourceActionsAreDatasourceReadAndWrite() + { + // READ and WRITE, matching the Coordinator catalog API's requirement for write operations: the write is forwarded + // with escalated credentials, so this statement-level check is the only one the issuing user faces. + final DirectStatement stmt = statement("CREATE TABLE tbl (a VARCHAR)"); + stmt.execute(); + assertEquals( + ImmutableSet.of( + new ResourceAction(new Resource("tbl", ResourceType.DATASOURCE), Action.READ), + new ResourceAction(new Resource("tbl", ResourceType.DATASOURCE), Action.WRITE) + ), + ImmutableSet.copyOf(stmt.resources()) + ); + } + + /** + * The grammar cannot express EXPLAIN of a DDL statement, so the planner's defensive guard is exercised directly: + * parse the CREATE, wrap it in a hand-built SqlExplain, and hand that to the planner as if the parser had produced + * it. + */ + @Test + public void testExplainOfDdlIsRejectedByThePlanner() + { + final SqlNode create = DruidSqlParser.parse("CREATE TABLE tbl (a VARCHAR)", true).getMainStatement(); + final SqlExplain explain = new SqlExplain( + SqlParserPos.ZERO, + create, + SqlExplainLevel.ALL_ATTRIBUTES.symbol(SqlParserPos.ZERO), + SqlExplain.Depth.PHYSICAL.symbol(SqlParserPos.ZERO), + SqlExplainFormat.TEXT.symbol(SqlParserPos.ZERO), + 0 + ); + final PlannerFactory plannerFactory = queryFramework() + .plannerFixture(PlannerConfig.builder().enableCatalogDdl(true).build(), new AuthConfig()) + .plannerFactory(); + try (DruidPlanner planner = plannerFactory.createPlanner( + queryFramework().engine(), + "EXPLAIN PLAN FOR CREATE TABLE tbl (a VARCHAR)", + explain, + CalciteTests.SUPER_USER_AUTH_RESULT, + Collections.emptySet(), + Collections.emptyMap(), + null + )) { + final DruidException e = assertThrows(DruidException.class, planner::validate); + assertTrue(e.getMessage().contains("EXPLAIN is not supported for [CREATE_TABLE]"), e.getMessage()); + } + assertTrue(WRITER.calls.isEmpty()); + } + + /** + * An unauthorized statement must not touch the catalog: authorization runs before planning, and the write itself is + * deferred to the result's run, so a rejected statement leaves no trace. Covers both denial modes: no permissions at + * all, and READ without WRITE (DDL requires both). + */ + @Test + public void testUnauthorizedDdlNeverReachesTheWriter() + { + assertThrows( + ForbiddenException.class, + () -> statement("CREATE TABLE forbiddenDatasource (a VARCHAR)", CalciteTests.REGULAR_USER_AUTH_RESULT) + .execute() + ); + assertThrows( + ForbiddenException.class, + () -> statement("CREATE TABLE readOnlyTbl (a VARCHAR)", CalciteTests.REGULAR_USER_AUTH_RESULT).execute() + ); + assertTrue(WRITER.calls.isEmpty()); + } + + /** + * Planning a DDL statement is free of side effects; the catalog write happens when the result runs. + */ + @Test + public void testWriteHappensOnRunNotOnPlan() + { + final DirectStatement.ResultSet resultSet = statement("CREATE TABLE tbl (a VARCHAR)").plan(); + assertTrue(WRITER.calls.isEmpty()); + + resultSet.run(); + assertEquals(1, WRITER.calls.size()); + assertEquals("createTable", WRITER.calls.get(0).operation); + } + + @Test + public void testDdlReturnsNoRows() + { + final DirectStatement stmt = statement("CREATE TABLE tbl (a VARCHAR)"); + final List<Object[]> results = stmt.execute().getResults().toList(); + assertEquals(ImmutableList.of(), results); + } + + @Test + public void testAlterTableAddColumn() + { + // ADD and ALTER differ only in which existence outcome is an error, and that rule is enforced inside the + // Coordinator's update transaction, so all the statement does is pick the verb. EditorTest covers the rules. + execute("ALTER TABLE tbl ADD COLUMN b BIGINT"); + + final RecordingCatalogTableWriter.Call call = WRITER.lastCall("addColumns"); + assertEquals(ImmutableList.of(new ColumnSpec("b", "BIGINT", null)), call.columns); + } + + @Test + public void testAlterTableDropColumn() + { + execute("ALTER TABLE tbl DROP COLUMN gone"); + assertEquals(ImmutableList.of("gone"), WRITER.lastCall("dropColumns").droppedColumns); + } + + @Test + public void testAlterTableAlterColumn() + { + execute("ALTER TABLE tbl ALTER COLUMN cnt SET DATA TYPE DOUBLE"); + assertEquals( + ImmutableList.of(new ColumnSpec("cnt", "DOUBLE", null)), + WRITER.lastCall("alterColumns").columns + ); + } + + @Test + public void testAlterTableSetProperties() + { + execute("ALTER TABLE tbl SET PROPERTIES (targetSegmentRows = 3000000, sealed = TRUE, description = 'hi')"); + + final Map<String, Object> properties = WRITER.lastCall("updateProperties").properties; + assertEquals(3000000L, properties.get("targetSegmentRows")); + assertEquals(true, properties.get("sealed")); + assertEquals("hi", properties.get("description")); + } + + @Test + public void testAlterTableSetPropertyToNullRemovesIt() + { + execute("ALTER TABLE tbl SET PROPERTIES (description = NULL)"); + final Map<String, Object> properties = WRITER.lastCall("updateProperties").properties; + assertTrue(properties.containsKey("description")); + assertNull(properties.get("description")); + } + + @Test + public void testCreateTableRejectsDuplicateColumn() + { + final DruidException e = assertThrows( + DruidException.class, + () -> execute("CREATE TABLE tbl (a VARCHAR, a BIGINT)") + ); + assertTrue(e.getMessage().contains("Column [a] is declared more than once")); + } + + @Test + public void testCreateTableRejectsUnsupportedType() + { + final DruidException e = assertThrows( + DruidException.class, + () -> execute("CREATE TABLE tbl (a TYPE('NOT_A_TYPE'))") + ); + assertTrue(e.getMessage().contains("unsupported type")); + } + + /** + * Any spelling that resolves to a LONG is accepted for the time column, and is stored as written. + */ + @Test + public void testCreateTableTimeColumnSpellings() + { + execute("CREATE TABLE tbl (__time BIGINT)"); + assertEquals(ImmutableList.of(new ColumnSpec("__time", "BIGINT", null)), WRITER.calls.get(0).spec.columns()); + + WRITER.reset(); + execute("CREATE TABLE tbl (__time TYPE('LONG'))"); + assertEquals(ImmutableList.of(new ColumnSpec("__time", "LONG", null)), WRITER.calls.get(0).spec.columns()); + } + + @Test + public void testCreateTableRejectsNonLongTimeColumn() + { + final DruidException e = assertThrows( + DruidException.class, + () -> execute("CREATE TABLE tbl (__time VARCHAR)") + ); + assertTrue(e.getMessage().contains("Column [__time] must have type")); + } + + @Test + public void testCreateTableRejectsNonDruidSchema() + { + final DruidException e = assertThrows( + DruidException.class, + () -> execute("CREATE TABLE lookup.tbl (a VARCHAR)") + ); + assertTrue(e.getMessage().contains("is not a Druid datasource")); + } + + @Test + public void testCreateTableRejectsBothReplaceAndIfNotExists() + { + final DruidException e = assertThrows( + DruidException.class, + () -> execute("CREATE OR REPLACE TABLE IF NOT EXISTS tbl (a VARCHAR)") + ); + assertTrue(e.getMessage().contains("Cannot specify both OR REPLACE and IF NOT EXISTS")); + } + + @Test + public void testCreateTableRejectsClusteringExpression() + { + final DruidException e = assertThrows( + DruidException.class, + () -> execute("CREATE TABLE tbl (a VARCHAR) CLUSTERED BY a DESC") + ); + assertTrue(e.getMessage().contains("must be a column name")); + } + + /** + * The feature is off unless an operator turns it on, so that upgrading a cluster does not silently widen what a + * datasource WRITE permission allows. + */ + @Test + public void testDdlIsDisabledByDefault() + { + final DirectStatement stmt = getSqlStatementFactory(PlannerConfig.builder().build(), new AuthConfig()) + .directStatement( + SqlQueryPlus.builder("CREATE TABLE tbl (a VARCHAR)") + .auth(CalciteTests.SUPER_USER_AUTH_RESULT) + .build() + ); + final DruidException e = assertThrows(DruidException.class, stmt::execute); + assertTrue(e.getMessage().contains("druid.sql.planner.enableCatalogDdl"), e.getMessage()); + assertEquals(ImmutableList.of(), WRITER.calls); + } + + /** + * The stored specification must be the one the planner would produce for the equivalent query, since that is what + * makes a projection match at query time. Pinned as JSON so a change in planner output is visible here. + */ + @Test + public void testCreateTableWithProjection() throws Exception + { + execute( + "CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT," + + " PROJECTION daily AS (SELECT TIME_FLOOR(__time, 'P1D'), page, SUM(cnt) AS total GROUP BY 1, 2))" + ); + + assertEquals( + "[{\"spec\":{\"type\":\"aggregate\",\"name\":\"daily\"," + + "\"virtualColumns\":[{\"type\":\"expression\",\"name\":\"v0\"," + + "\"expression\":\"timestamp_floor(\\\"__time\\\",'P1D',null,'UTC')\",\"outputType\":\"LONG\"}]," + + "\"groupingColumns\":[{\"type\":\"long\",\"name\":\"v0\",\"multiValueHandling\":\"SORTED_ARRAY\"," + + "\"createBitmapIndex\":false},{\"type\":\"string\",\"name\":\"page\"," + + "\"multiValueHandling\":\"SORTED_ARRAY\",\"createBitmapIndex\":true}]," + + "\"aggregators\":[{\"type\":\"longSum\",\"name\":\"total\",\"fieldName\":\"cnt\"}]," + + "\"ordering\":[{\"columnName\":\"v0\",\"order\":\"ascending\"}," + + "{\"columnName\":\"page\",\"order\":\"ascending\"}]}}]", + projectionsJson() + ); + } + + /** + * A projection body is planned under the statement's own context, so a SET clause that changes how the equivalent + * query would plan changes the stored definition the same way. Here the session time zone reaches the TIME_FLOOR. + */ + @Test + public void testProjectionBodyHonorsStatementContext() throws Exception + { + execute( + "SET sqlTimeZone = 'America/Los_Angeles';\n" + + "CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT," + + " PROJECTION daily AS (SELECT TIME_FLOOR(__time, 'P1D'), page, SUM(cnt) AS total GROUP BY 1, 2))" + ); + + assertTrue( + projectionsJson().contains("timestamp_floor(\\\"__time\\\",'P1D',null,'America/Los_Angeles')"), + projectionsJson() + ); + } + + /** + * The overrides the lift depends on are applied on top of the statement's context, so a SET clause cannot put the + * planner into a shape the lift does not understand. + */ + @Test + public void testProjectionBodyContextCannotOverrideDeterministicOverrides() throws Exception + { + execute( + "SET sqlUseGranularity = TRUE;\n" + + "CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT," + + " PROJECTION daily AS (SELECT TIME_FLOOR(__time, 'P1D'), page, SUM(cnt) AS total GROUP BY 1, 2))" + ); + + // Still lifted as an ordinary grouping column rather than a query granularity, exactly as without the SET. + assertTrue( + projectionsJson().contains("timestamp_floor(\\\"__time\\\",'P1D',null,'UTC')"), + projectionsJson() + ); + } + + /** + * A projection defined with TIME_FLOOR must carry a granularity the segment layer can recover, which is how the + * projection gets matched to time-grouped queries. + */ + @Test + public void testProjectionGranularityIsRecoverable() + { + execute( + "CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT," + + " PROJECTION hourly AS (SELECT TIME_FLOOR(__time, 'PT1H'), page, SUM(cnt) AS total GROUP BY 1, 2))" + ); + + final AggregateProjectionSpec spec = projection(0).getSpec(); + final String timeColumn = spec.toMetadataSchema().getTimeColumnName(); + assertEquals("v0", timeColumn); + assertEquals( + Granularities.HOUR, + Granularities.fromVirtualColumn(spec.getVirtualColumns().getVirtualColumn(timeColumn)) + ); + } + + @Test + public void testProjectionWithFilter() + { + execute( + "CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT," + + " PROJECTION filtered AS (SELECT page, SUM(cnt) AS total WHERE page <> 'skip' GROUP BY page))" + ); + assertEquals("!page = skip", projection(0).getSpec().getFilter().toString()); + } + + /** + * A time bound written in the body is moved into the query's intervals during planning, and has to be put back: + * a projection stores a filter, not an interval. + */ + @Test + public void testProjectionWithTimeFilter() + { + execute( + "CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT," + + " PROJECTION recent AS (SELECT page, SUM(cnt) AS total" + + " WHERE __time >= TIMESTAMP '2020-01-01 00:00:00' GROUP BY page))" + ); + assertNotNull(projection(0).getSpec().getFilter(), "time filter must survive as a filter"); Review Comment: Good eye, there is not currently logic that can handle matching __time filters in projections against query intervals, so this is currently modeling building an unreachable projection. Claude came up with this, I left it because i think retaining them as filters makes more sense than trying to give a standalone interval to projections, so i view at as the right form and since this test is just covering projection construction and not query time matching it seemed harmless to leave. That is, assuming I add support for special handling projections with __time filters that are comparable to intervals as a follow-up, which I think it would be pretty straightforward to add, so I have it on my list. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
