gianm commented on code in PR #19830:
URL: https://github.com/apache/druid/pull/19830#discussion_r3994879403


##########
docs/development/extensions-core/catalog.md:
##########
@@ -43,6 +43,208 @@ allowing queries to be more concise, and simpler to write. 
This also allows the
 written into a defined column of the table is consistent with that columns 
definition, minimizing errors where unexpected
 data is written into a particular column of the table.
 
+### Effects of a table definition
+
+A table definition takes effect when it is read, not when it is written: 
changing one, whether through SQL DDL or the
+REST API, rewrites nothing. What it affects:
+
+- **SELECT queries** report each declared column with its declared type, in 
declared order. Segments that store a
+  column as a different physical type are converted to the declared type at 
query time. Columns present in segments
+  but not declared in the catalog remain queryable, and follow the declared 
columns with their physical types.
+- **INSERT and REPLACE** (the catalog applies only to SQL-based ingestion) 
validate against the definition rather
+  than merely defaulting from it. A query column targeting a declared column 
must be convertible to the declared type
+  without changing it: inserting a `BIGINT` value into a column declared 
`VARCHAR` stores the value as a string, while
+  inserting a `VARCHAR` value into a column declared `BIGINT` is an error. 
Columns the query produces that the table
+  does not declare are rejected when the table is 
[`sealed`](#table-properties) and ingested normally otherwise.
+  Table properties such as `segmentGranularity` and `clusterKeys` act as 
defaults that an individual statement may
+  override, such as with its own `PARTITIONED BY`.
+- **Streaming and native batch ingestion** do not consult the catalog; their 
specs are unaffected by any table
+  definition.
+
+### SQL DDL
+
+Tables can be defined with SQL instead of by posting a table specification. 
`CREATE TABLE` and `ALTER TABLE` are
+submitted to the Broker like any other SQL statement, and write the same 
catalog metadata the REST API does. They
+return no rows.
+
+These statements change catalog metadata only. They never create, modify, or 
delete segments: defining a table does
+not ingest anything, and altering a column does not rewrite existing data. 
Changes take effect for subsequent queries
+and ingestion, as described in [Effects of a table 
definition](#effects-of-a-table-definition).
+
+These statements are disabled by default. Set 
`druid.sql.planner.enableCatalogDdl` to `true` on the Broker to enable
+them. They require both `READ` and `WRITE` permission on the datasource, the 
same permissions the catalog API
+requires for its write operations, so enabling them lets anyone who can ingest 
into a datasource also change its
+catalog definition; leave them disabled if you manage catalog entries with 
your own tooling. The setting cannot be
+overridden per query.
+
+The `druid-catalog` extension must be loaded on both the Broker and the 
Coordinator; without it, these statements
+report that the extension is not available.
+
+```sql
+CREATE [OR REPLACE] TABLE [IF NOT EXISTS] <table>
+  [ [ SEALED ] ( <table element> [, ...] ) ]
+  [ PARTITIONED BY <granularity> ]
+  [ CLUSTERED BY <column> [, ...] ]
+
+<table element> ::=
+    <column> <type>
+  | PROJECTION <name> AS ( <select> )
+```
+
+`OR REPLACE` replaces the specification of an existing table; `IF NOT EXISTS` 
leaves an existing table unchanged.
+The two cannot be combined. `PARTITIONED BY` sets 
[`segmentGranularity`](#table-properties) and `CLUSTERED BY` sets
+`clusterKeys`, both of which a later `INSERT` or `REPLACE` inherits unless it 
states its own. `SEALED` sets
+[`sealed`](#table-properties), which requires every ingested column to be 
declared. It is written just before the
+column list because that is what it describes: the list is the table's whole 
schema, so the statement does not accept
+`SEALED` without one. (The `sealed` property itself can still be set on any 
table through `SET PROPERTIES` or the
+REST API.)
+
+Note that the table-level `CLUSTERED BY` is a sort order applied to each 
ingestion, which is a different thing from
+the `CLUSTERED BY` inside a [`__base` projection](#the-base-table), which 
defines how segments physically group rows.
+
+Column types are written as SQL types, such as `VARCHAR`, `BIGINT`, `DOUBLE`, 
or `VARCHAR ARRAY`. The `__time` column
+is written as `TIMESTAMP`. Types that have no SQL spelling, such as complex 
types, use `TYPE('...')` with the Druid
+native type string:
+
+```sql
+CREATE TABLE "druid"."visits" (
+  __time TIMESTAMP,
+  user_id VARCHAR,
+  pages_visited BIGINT,
+  sketch TYPE('COMPLEX<thetaSketch>')
+)
+PARTITIONED BY DAY
+CLUSTERED BY user_id
+```
+
+`ALTER TABLE` supports one change per statement, so that each statement is a 
single atomic catalog operation:
+
+```sql
+ALTER TABLE <table> ADD COLUMN <column> <type>
+ALTER TABLE <table> DROP COLUMN <column>
+ALTER TABLE <table> ALTER COLUMN <column> SET DATA TYPE <type>
+ALTER TABLE <table> ADD [IF NOT EXISTS] PROJECTION <name> AS ( <select> )

Review Comment:
   Please change it to `ADD PROJECTION [IF NOT EXISTS] <name>`, as this is more 
commonly the way such syntaxes are done in SQL systems. It also better matches 
`DROP PROJECTION [IF EXISTS]`.



##########
sql/src/main/java/org/apache/druid/sql/calcite/planner/CatalogColumnTypes.java:
##########
@@ -0,0 +1,145 @@
+/*
+ * 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.planner;
+
+import org.apache.calcite.avatica.SqlType;
+import org.apache.calcite.sql.SqlDataTypeSpec;
+import org.apache.calcite.sql.SqlIdentifier;
+import org.apache.calcite.sql.SqlTypeNameSpec;
+import org.apache.calcite.sql.SqlUserDefinedTypeNameSpec;
+import org.apache.calcite.sql.type.SqlTypeName;
+import org.apache.druid.catalog.model.Columns;
+import org.apache.druid.java.util.common.IAE;
+import org.apache.druid.java.util.common.StringUtils;
+import org.apache.druid.segment.column.ColumnType;
+
+/**
+ * Converts a SQL type as written in a statement into the type string stored 
in the catalog. Calcite accepts any
+ * type spelling; Druid accepts only its own supported types and their 
aliases, plus the {@code TYPE('...')} escape
+ * hatch for native type strings that have no SQL spelling, such as {@code 
COMPLEX<json>}.
+ * <p>
+ * Druid has its own rules for nullability, so any nullability clause is 
ignored.
+ */
+public class CatalogColumnTypes
+{
+  /**
+   * Which statement the type was written in. The two differ in what they 
accept, so they are named rather than
+   * flagged.
+   */
+  private enum Target
+  {
+    /**
+     * The {@code EXTEND} clause, describing columns read from an external 
input source. Those are read as their
+     * underlying storage type, so {@code TIMESTAMP} is not among them, and 
the escape hatch is limited to complex
+     * types.
+     */
+    EXTERNAL,
+
+    /**
+     * A catalog DDL statement, which additionally accepts {@code TIMESTAMP} 
(how {@code __time} is spelled in SQL)
+     * and any native type string through the escape hatch.
+     */
+    CATALOG
+  }
+
+  private CatalogColumnTypes()
+  {
+    // No instantiation.
+  }
+
+  public static String forExternalColumn(String name, SqlDataTypeSpec dataType)
+  {
+    return convert(name, dataType, Target.EXTERNAL);
+  }
+
+  public static String forCatalogColumn(String name, SqlDataTypeSpec dataType)
+  {
+    final String typeString = convert(name, dataType, Target.CATALOG);
+    // the catalog rejects unparseable types at write time, but catching it 
here attributes the error to the statement
+    // rather than to a Coordinator round trip.
+    if (Columns.druidTypeFromString(typeString) == null) {
+      throw unsupportedType(name, dataType);
+    }
+    return typeString;
+  }
+
+  private static String convert(String name, SqlDataTypeSpec dataType, Target 
target)
+  {
+    final SqlTypeNameSpec spec = dataType.getTypeNameSpec();
+    if (spec == null) {
+      throw unsupportedType(name, dataType);
+    }
+    final SqlIdentifier typeNameIdentifier = spec.getTypeName();
+    if (typeNameIdentifier == null || !typeNameIdentifier.isSimple()) {
+      throw unsupportedType(name, dataType);
+    }
+    final String simpleName = typeNameIdentifier.getSimple();
+
+    if (spec instanceof SqlUserDefinedTypeNameSpec) {
+      // The TYPE('...') escape hatch names a Druid native type. Parse and 
validate rather than passing the raw
+      // string downstream, where a malformed type string would silently 
resolve to a different type, and return the
+      // canonical form.
+      if (target == Target.EXTERNAL && 
!StringUtils.toLowerCase(simpleName).startsWith("complex<")) {

Review Comment:
   Any reason to not allow non-complex types here? Seems like it would be fine 
to allow arrays and primitives.



##########
sql/src/main/java/org/apache/druid/sql/calcite/planner/CatalogColumnTypes.java:
##########
@@ -0,0 +1,145 @@
+/*
+ * 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.planner;
+
+import org.apache.calcite.avatica.SqlType;
+import org.apache.calcite.sql.SqlDataTypeSpec;
+import org.apache.calcite.sql.SqlIdentifier;
+import org.apache.calcite.sql.SqlTypeNameSpec;
+import org.apache.calcite.sql.SqlUserDefinedTypeNameSpec;
+import org.apache.calcite.sql.type.SqlTypeName;
+import org.apache.druid.catalog.model.Columns;
+import org.apache.druid.java.util.common.IAE;
+import org.apache.druid.java.util.common.StringUtils;
+import org.apache.druid.segment.column.ColumnType;
+
+/**
+ * Converts a SQL type as written in a statement into the type string stored 
in the catalog. Calcite accepts any
+ * type spelling; Druid accepts only its own supported types and their 
aliases, plus the {@code TYPE('...')} escape
+ * hatch for native type strings that have no SQL spelling, such as {@code 
COMPLEX<json>}.
+ * <p>
+ * Druid has its own rules for nullability, so any nullability clause is 
ignored.
+ */
+public class CatalogColumnTypes
+{
+  /**
+   * Which statement the type was written in. The two differ in what they 
accept, so they are named rather than
+   * flagged.
+   */
+  private enum Target
+  {
+    /**
+     * The {@code EXTEND} clause, describing columns read from an external 
input source. Those are read as their
+     * underlying storage type, so {@code TIMESTAMP} is not among them, and 
the escape hatch is limited to complex
+     * types.
+     */
+    EXTERNAL,
+
+    /**
+     * A catalog DDL statement, which additionally accepts {@code TIMESTAMP} 
(how {@code __time} is spelled in SQL)
+     * and any native type string through the escape hatch.
+     */
+    CATALOG
+  }
+
+  private CatalogColumnTypes()
+  {
+    // No instantiation.
+  }
+
+  public static String forExternalColumn(String name, SqlDataTypeSpec dataType)
+  {
+    return convert(name, dataType, Target.EXTERNAL);
+  }
+
+  public static String forCatalogColumn(String name, SqlDataTypeSpec dataType)
+  {
+    final String typeString = convert(name, dataType, Target.CATALOG);
+    // the catalog rejects unparseable types at write time, but catching it 
here attributes the error to the statement
+    // rather than to a Coordinator round trip.
+    if (Columns.druidTypeFromString(typeString) == null) {
+      throw unsupportedType(name, dataType);
+    }
+    return typeString;
+  }
+
+  private static String convert(String name, SqlDataTypeSpec dataType, Target 
target)
+  {
+    final SqlTypeNameSpec spec = dataType.getTypeNameSpec();
+    if (spec == null) {
+      throw unsupportedType(name, dataType);
+    }
+    final SqlIdentifier typeNameIdentifier = spec.getTypeName();
+    if (typeNameIdentifier == null || !typeNameIdentifier.isSimple()) {
+      throw unsupportedType(name, dataType);
+    }
+    final String simpleName = typeNameIdentifier.getSimple();
+
+    if (spec instanceof SqlUserDefinedTypeNameSpec) {
+      // The TYPE('...') escape hatch names a Druid native type. Parse and 
validate rather than passing the raw
+      // string downstream, where a malformed type string would silently 
resolve to a different type, and return the
+      // canonical form.
+      if (target == Target.EXTERNAL && 
!StringUtils.toLowerCase(simpleName).startsWith("complex<")) {
+        throw unsupportedType(name, dataType);
+      }
+      final ColumnType nativeType = ColumnType.fromString(simpleName);
+      if (nativeType == null) {

Review Comment:
   Should check that the complex type is valid, or else this will cause 
failures later when we attempt to apply the table spec, leading to a broken 
table.



##########
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");
+    
assertTrue(projection(0).getSpec().getFilter().getRequiredColumns().contains("__time"));
+  }
+
+  @Test
+  public void testProjectionSelectDistinct()
+  {
+    execute("CREATE TABLE tbl (a VARCHAR, PROJECTION d AS (SELECT DISTINCT 
a))");
+    final AggregateProjectionSpec spec = projection(0).getSpec();
+    assertEquals(1, spec.getGroupingColumns().size());
+    assertEquals("a", spec.getGroupingColumns().get(0).getName());
+    assertEquals(0, spec.getAggregators().length);
+  }
+
+  @Test
+  public void testMultipleProjections()
+  {
+    execute(
+        "CREATE TABLE tbl (a VARCHAR, b BIGINT,"
+        + " PROJECTION p1 AS (SELECT a, SUM(b) AS s GROUP BY a),"
+        + " PROJECTION p2 AS (SELECT b, COUNT(*) AS c GROUP BY b))"
+    );
+    assertEquals(List.of("p1", "p2"), 
List.of(projection(0).getSpec().getName(), projection(1).getSpec().getName()));
+  }
+
+  @Test
+  public void testAlterTableAddProjection()
+  {
+    WRITER.existing.put(TableId.datasource("tbl"), tableWithColumns("a"));
+    execute("ALTER TABLE tbl ADD PROJECTION p AS (SELECT a, COUNT(*) AS c 
GROUP BY a)");
+
+    final RecordingCatalogTableWriter.Call call = 
WRITER.lastCall("addProjection");
+    assertEquals("p", call.projection.getSpec().getName());
+    assertFalse(call.ifNotExists);
+  }
+
+  @Test
+  public void testAlterTableAddProjectionIfNotExists()
+  {
+    WRITER.existing.put(TableId.datasource("tbl"), tableWithColumns("a"));
+    execute("ALTER TABLE tbl ADD IF NOT EXISTS PROJECTION p AS (SELECT a GROUP 
BY a)");
+    assertTrue(WRITER.lastCall("addProjection").ifNotExists);
+  }
+
+  @Test
+  public void testAlterTableDropProjection()
+  {
+    execute("ALTER TABLE tbl DROP PROJECTION p");
+    final RecordingCatalogTableWriter.Call call = 
WRITER.lastCall("dropProjection");
+    assertEquals("p", call.projectionName);
+    assertFalse(call.ifExists);
+
+    WRITER.reset();
+    execute("ALTER TABLE tbl DROP PROJECTION IF EXISTS p");
+    assertTrue(WRITER.lastCall("dropProjection").ifExists);
+  }
+
+  @Test
+  public void testProjectionRejectsUnaliasedAggregate()
+  {
+    final DruidException e = assertThrows(
+        DruidException.class,
+        () -> execute("CREATE TABLE tbl (a VARCHAR, b BIGINT, PROJECTION p AS 
(SELECT a, SUM(b) GROUP BY a))")
+    );
+    assertTrue(e.getMessage().contains("no name"), e.getMessage());
+  }
+
+  @Test
+  public void testProjectionRejectsPostAggregation()
+  {
+    final DruidException e = assertThrows(
+        DruidException.class,
+        () -> execute("CREATE TABLE tbl (a VARCHAR, b BIGINT, PROJECTION p AS 
(SELECT a, AVG(b) AS m GROUP BY a))")
+    );
+    assertTrue(e.getMessage().contains("expression over aggregates"), 
e.getMessage());
+  }
+
+  /**
+   * Subqueries survive the grammar (they are just expressions in the select 
list or WHERE), so the translator walks
+   * the body and rejects them with a pointed message rather than letting the 
planner produce something the lift
+   * cannot express.
+   */
+  @Test
+  public void testProjectionRejectsSubquery()
+  {
+    for (String body : new String[]{
+        "SELECT a, (SELECT 1) AS s GROUP BY 1, 2",
+        "SELECT a WHERE a IN (SELECT a) GROUP BY a",
+    }) {
+      final DruidException e = assertThrows(
+          DruidException.class,
+          () -> execute("CREATE TABLE tbl (a VARCHAR, PROJECTION p AS (" + 
body + "))"),
+          body
+      );
+      assertTrue(e.getMessage().contains("contains a subquery"), body + " -> " 
+ e.getMessage());
+    }
+  }
+
+  /**
+   * Window functions survive the grammar as ordinary select items, and are 
rejected by the nested planner because the
+   * projection engine does not declare {@code 
EngineFeature.WINDOW_FUNCTIONS}; the error is attributed to the
+   * projection by name.
+   */
+  @Test
+  public void testProjectionRejectsWindowFunctions()
+  {
+    final DruidException e = assertThrows(
+        DruidException.class,
+        () -> execute(
+            "CREATE TABLE tbl (a VARCHAR, b BIGINT,"
+            + " PROJECTION p AS (SELECT a, SUM(b) OVER (PARTITION BY a) AS w 
GROUP BY a, b))"
+        )
+    );
+    assertTrue(e.getMessage().contains("Cannot define projection [p]"), 
e.getMessage());
+    assertTrue(e.getMessage().contains("window functions"), e.getMessage());
+  }
+
+  @Test
+  public void testProjectionRejectsUnknownColumn()
+  {
+    final DruidException e = assertThrows(
+        DruidException.class,
+        () -> execute("CREATE TABLE tbl (a VARCHAR, PROJECTION p AS (SELECT 
nope GROUP BY nope))")
+    );
+    assertTrue(e.getMessage().contains("nope"), e.getMessage());
+  }
+
+  @Test
+  public void testProjectionRejectsNonAggregatingBody()
+  {
+    final DruidException e = assertThrows(
+        DruidException.class,
+        () -> execute("CREATE TABLE tbl (a VARCHAR, PROJECTION p AS (SELECT 
a))")
+    );
+    assertTrue(e.getMessage().contains("does not aggregate"), e.getMessage());
+  }
+
+  /**
+   * {@code __base} names the table's own layout and is handled separately; 
every other name beginning with the
+   * reserved prefix stays unavailable.
+   */
+  @Test
+  public void testProjectionRejectsOtherReservedNames()
+  {
+    final DruidException e = assertThrows(
+        DruidException.class,
+        () -> execute("CREATE TABLE tbl (a VARCHAR, PROJECTION __other AS 
(SELECT a GROUP BY a))")
+    );
+    assertTrue(e.getMessage().contains("reserved name"), e.getMessage());
+  }
+
+  @Test
+  public void testProjectionRejectsDuplicateName()
+  {
+    final DruidException e = assertThrows(
+        DruidException.class,
+        () -> execute(
+            "CREATE TABLE tbl (a VARCHAR, PROJECTION p AS (SELECT a GROUP BY 
a),"
+            + " PROJECTION p AS (SELECT a GROUP BY a))"
+        )
+    );
+    assertTrue(e.getMessage().contains("declared more than once"), 
e.getMessage());
+  }
+
+  /**
+   * The reserved {@code __base} projection describes the table's own layout, 
so it becomes the baseTable property
+   * rather than one of the projections. A computed column becomes a virtual 
column materializing the declared column
+   * it fills.
+   */
+  @Test
+  public void testCreateTableWithBaseProjection() throws Exception
+  {
+    execute(
+        "CREATE TABLE tbl SEALED ("
+        + " tenant VARCHAR,"
+        + " bucket BIGINT,"
+        + " __time TIMESTAMP,"
+        + " user_id BIGINT,"
+        + " PROJECTION __base AS ("
+        + "   SELECT tenant, ABS(user_id) AS bucket, __time, user_id"
+        + "   CLUSTERED BY tenant, bucket"
+        + " )"
+        + ") PARTITIONED BY DAY"
+    );
+
+    final TableSpec spec = WRITER.calls.get(0).spec;
+    assertEquals(true, spec.properties().get(DatasourceDefn.SEALED_PROPERTY));
+    
assertNull(spec.properties().get(DatasourceDefn.PROJECTIONS_KEYS_PROPERTY));
+    assertEquals(
+        "{\"clusteringColumns\":[\"tenant\",\"bucket\"],"
+        + "\"virtualColumns\":[{\"type\":\"expression\",\"name\":\"bucket\","
+        + "\"expression\":\"abs(\\\"user_id\\\")\",\"outputType\":\"LONG\"}],"
+        + "\"type\":\"clusteredValueGroups\"}",
+        queryFramework().queryJsonMapper()
+                        
.writeValueAsString(spec.properties().get(DatasourceDefn.BASE_TABLE_PROPERTY))
+    );
+  }
+
+  /**
+   * A base table column is stored under the name it declares, so a body item 
that renames another column is rejected:
+   * only a virtual column materializes a name the body did not select, and a 
bare reference produces none.
+   * <p>
+   * Selecting the same column twice also makes the scan's deduplicated column 
list shorter than the select list, so
+   * the source of each item is taken from the query's output signature. 
Reading the scan's columns positionally used
+   * to run off the end of that list and fail as an internal error.
+   */
+  @Test
+  public void testBaseProjectionRenamedColumnFails()
+  {
+    final DruidException e = assertThrows(
+        DruidException.class,
+        () -> execute(
+            "CREATE TABLE tbl SEALED (id BIGINT, copy BIGINT, __time 
TIMESTAMP,"
+            + " PROJECTION __base AS (SELECT id, id AS copy, __time CLUSTERED 
BY id))"
+        )
+    );
+    assertTrue(
+        e.getMessage().contains("its column 2 selects [id] but declares it as 
[copy]"),
+        e.getMessage()
+    );
+  }
+
+  @Test
+  public void testBaseProjectionWithoutComputedColumns()
+  {
+    execute(
+        "CREATE TABLE tbl SEALED (tenant VARCHAR, __time TIMESTAMP, v BIGINT,"
+        + " PROJECTION __base AS (SELECT tenant, __time, v CLUSTERED BY 
tenant))"
+    );
+    final ClusteredValueGroupsBaseTableMetadata baseTable = baseTable();
+    assertEquals(List.of("tenant"), baseTable.getClusteringColumns());
+    assertEquals(0, baseTable.getVirtualColumns().getVirtualColumns().length);
+  }
+
+  /**
+   * A base table and aggregate projections are different catalog entities and 
may coexist.
+   */
+  @Test
+  public void testBaseProjectionAlongsideAggregateProjection()
+  {
+    execute(
+        "CREATE TABLE tbl SEALED (tenant VARCHAR, __time TIMESTAMP, v BIGINT,"
+        + " PROJECTION __base AS (SELECT tenant, __time, v CLUSTERED BY 
tenant),"
+        + " PROJECTION by_tenant AS (SELECT tenant, SUM(v) AS sum_v GROUP BY 
tenant))"
+    );
+    final TableSpec spec = WRITER.calls.get(0).spec;
+    assertNotNull(spec.properties().get(DatasourceDefn.BASE_TABLE_PROPERTY));

Review Comment:
   This test class has a lot of `assertNotNull` that really should be asserting 
equality to the expected thing. Please tighten them up.



##########
sql/src/main/java/org/apache/druid/sql/calcite/planner/CatalogDdlHandler.java:
##########
@@ -0,0 +1,702 @@
+/*
+ * 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.planner;
+
+import com.google.common.collect.ImmutableList;
+import com.google.common.collect.ImmutableSet;
+import com.google.common.collect.Iterables;
+import org.apache.calcite.jdbc.CalciteSchema;
+import org.apache.calcite.rel.type.RelDataType;
+import org.apache.calcite.rel.type.RelDataTypeFactory;
+import org.apache.calcite.sql.SqlIdentifier;
+import org.apache.calcite.sql.SqlLiteral;
+import org.apache.calcite.sql.SqlNode;
+import org.apache.calcite.sql.SqlNodeList;
+import org.apache.calcite.sql.type.SqlTypeName;
+import org.apache.druid.catalog.model.ClusteredValueGroupsBaseTableMetadata;
+import org.apache.druid.catalog.model.ColumnSpec;
+import org.apache.druid.catalog.model.Columns;
+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.common.utils.IdUtils;
+import org.apache.druid.error.DruidException;
+import org.apache.druid.error.InvalidSqlInput;
+import org.apache.druid.java.util.common.IAE;
+import org.apache.druid.java.util.common.granularity.Granularities;
+import org.apache.druid.java.util.common.granularity.Granularity;
+import org.apache.druid.java.util.common.granularity.PeriodGranularity;
+import org.apache.druid.java.util.common.guava.Sequences;
+import org.apache.druid.query.explain.ExplainAttributes;
+import org.apache.druid.segment.column.ColumnType;
+import org.apache.druid.segment.projections.Projections;
+import org.apache.druid.server.QueryResponse;
+import org.apache.druid.server.security.Action;
+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.calcite.parser.DruidSqlAlterTable;
+import org.apache.druid.sql.calcite.parser.DruidSqlColumnDeclaration;
+import org.apache.druid.sql.calcite.parser.DruidSqlCreateTable;
+import org.apache.druid.sql.calcite.parser.DruidSqlParser;
+import org.apache.druid.sql.calcite.parser.DruidSqlPropertyAssignment;
+import org.apache.druid.sql.calcite.parser.SqlGranularityLiteral;
+import org.apache.druid.sql.calcite.parser.SqlProjectionSpec;
+import org.apache.druid.sql.calcite.run.EngineFeature;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+
+/**
+ * Handles the catalog DDL statements: {@code CREATE TABLE} and {@code ALTER 
TABLE}.
+ * <p>
+ * These statements are metadata operations, not queries. They are validated 
here, converted to a catalog
+ * {@link TableSpec} or column/property edit, and applied through {@link 
CatalogTableWriter}, which forwards them to
+ * the Coordinator. No Calcite validation or query planning takes place, and 
no rows are returned.
+ * <p>
+ * Validation is deliberately split. This class checks what it can attribute 
to a position in the statement (type
+ * spellings, duplicate columns, granularity, clustering) so that the error 
names the offending SQL. The Coordinator
+ * remains authoritative: {@code DatasourceDefn.validate} runs on write, and 
its message is surfaced verbatim.
+ */
+public abstract class CatalogDdlHandler extends 
SqlStatementHandler.BaseStatementHandler
+{
+  /**
+   * DDL produces no rows. A single-column type still has to be declared, 
because JDBC clients ask for a result set
+   * signature when preparing the statement.
+   */
+  private static final RelDataType RESULT_TYPE = resultType();
+
+  protected final SqlIdentifier tableIdentifier;
+  protected TableId tableId;
+
+  protected CatalogDdlHandler(SqlStatementHandler.HandlerContext 
handlerContext, SqlIdentifier tableIdentifier)
+  {
+    super(handlerContext);
+    this.tableIdentifier = tableIdentifier;
+  }
+
+  /**
+   * The runtime property that gates these statements. Read from {@link 
PlannerConfig} rather than the query context
+   * so that a user cannot turn the feature on for their own statement.
+   */
+  public static final String ENABLE_CATALOG_DDL_PROPERTY = 
"druid.sql.planner.enableCatalogDdl";
+
+  /**
+   * The reserved name of the base-table projection, which describes the 
physical layout of the table itself. Handled
+   * as a separate catalog property, not as one of the aggregate projections.
+   */
+  public static final String BASE_PROJECTION_NAME = "__base";

Review Comment:
   Use `Projections.BASE_TABLE_PROJECTION_NAME`?



##########
server/src/main/java/org/apache/druid/catalog/model/table/DatasourceDefn.java:
##########
@@ -110,6 +124,115 @@ public void validate(ResolvedTable table)
       // fail fast instead of surfacing layout problems at ingest time.
       baseTable.createSpec(table.spec().columns());
     }
+    validateProjections(table);
+  }
+
+  /**
+   * Cross-validate the declared projections. Names must be unique, a 
projection must not be coarser than the segments
+   * it lives in, and the types it groups by must agree with the types the 
table declares. For a sealed table the
+   * declared columns are the whole schema, so a projection that reads a 
column the table does not declare can never be
+   * built and is rejected; for a non-sealed table ingestion may add columns 
the catalog has not seen, so only the
+   * projections' internal consistency is checked.
+   */
+  private void validateProjections(ResolvedTable table)
+  {
+    final List<DatasourceProjectionMetadata> projections = 
table.decodeProperty(PROJECTIONS_KEYS_PROPERTY);
+    if (projections == null || projections.isEmpty()) {
+      return;
+    }
+
+    final List<AggregateProjectionSpec> specs = new 
ArrayList<>(projections.size());
+    for (DatasourceProjectionMetadata projection : projections) {
+      if (projection == null || projection.getSpec() == null) {
+        throw InvalidInput.exception("Projections must each have a [spec]");
+      }
+      specs.add(projection.getSpec());
+    }
+
+    final String granularity = 
table.stringProperty(SEGMENT_GRANULARITY_PROPERTY);
+    DataSchema.validateProjections(
+        specs,
+        granularity == null ? null : 
CatalogUtils.asDruidGranularity(granularity)
+    );
+
+    validateProjectionGroupingTypes(table, specs);
+
+    if (!table.booleanProperty(SEALED_PROPERTY) || table.spec().columns() == 
null) {
+      return;
+    }
+    final Set<String> declared = new 
HashSet<>(CatalogUtils.columnNames(table.spec().columns()));
+    declared.add(Columns.TIME_COLUMN);
+    for (AggregateProjectionSpec spec : specs) {
+      final Set<String> available = new HashSet<>(declared);
+      for (VirtualColumn virtualColumn : 
spec.getVirtualColumns().getVirtualColumns()) {
+        available.add(virtualColumn.getOutputName());
+      }
+      for (String required : requiredColumns(spec)) {
+        if (!available.contains(required)) {
+          throw InvalidInput.exception(
+              "Projection [%s] references column [%s], which table [%s] does 
not declare",
+              spec.getName(),
+              required,
+              table.spec().type()

Review Comment:
   Should be table name, not type, in this position.



##########
docs/development/extensions-core/catalog.md:
##########
@@ -43,6 +43,208 @@ allowing queries to be more concise, and simpler to write. 
This also allows the
 written into a defined column of the table is consistent with that columns 
definition, minimizing errors where unexpected
 data is written into a particular column of the table.
 
+### Effects of a table definition
+
+A table definition takes effect when it is read, not when it is written: 
changing one, whether through SQL DDL or the
+REST API, rewrites nothing. What it affects:
+
+- **SELECT queries** report each declared column with its declared type, in 
declared order. Segments that store a
+  column as a different physical type are converted to the declared type at 
query time. Columns present in segments
+  but not declared in the catalog remain queryable, and follow the declared 
columns with their physical types.
+- **INSERT and REPLACE** (the catalog applies only to SQL-based ingestion) 
validate against the definition rather
+  than merely defaulting from it. A query column targeting a declared column 
must be convertible to the declared type
+  without changing it: inserting a `BIGINT` value into a column declared 
`VARCHAR` stores the value as a string, while
+  inserting a `VARCHAR` value into a column declared `BIGINT` is an error. 
Columns the query produces that the table
+  does not declare are rejected when the table is 
[`sealed`](#table-properties) and ingested normally otherwise.
+  Table properties such as `segmentGranularity` and `clusterKeys` act as 
defaults that an individual statement may
+  override, such as with its own `PARTITIONED BY`.
+- **Streaming and native batch ingestion** do not consult the catalog; their 
specs are unaffected by any table
+  definition.
+
+### SQL DDL
+
+Tables can be defined with SQL instead of by posting a table specification. 
`CREATE TABLE` and `ALTER TABLE` are
+submitted to the Broker like any other SQL statement, and write the same 
catalog metadata the REST API does. They
+return no rows.
+
+These statements change catalog metadata only. They never create, modify, or 
delete segments: defining a table does
+not ingest anything, and altering a column does not rewrite existing data. 
Changes take effect for subsequent queries
+and ingestion, as described in [Effects of a table 
definition](#effects-of-a-table-definition).
+
+These statements are disabled by default. Set 
`druid.sql.planner.enableCatalogDdl` to `true` on the Broker to enable
+them. They require both `READ` and `WRITE` permission on the datasource, the 
same permissions the catalog API
+requires for its write operations, so enabling them lets anyone who can ingest 
into a datasource also change its
+catalog definition; leave them disabled if you manage catalog entries with 
your own tooling. The setting cannot be
+overridden per query.
+
+The `druid-catalog` extension must be loaded on both the Broker and the 
Coordinator; without it, these statements
+report that the extension is not available.
+
+```sql
+CREATE [OR REPLACE] TABLE [IF NOT EXISTS] <table>
+  [ [ SEALED ] ( <table element> [, ...] ) ]
+  [ PARTITIONED BY <granularity> ]
+  [ CLUSTERED BY <column> [, ...] ]
+
+<table element> ::=
+    <column> <type>
+  | PROJECTION <name> AS ( <select> )
+```
+
+`OR REPLACE` replaces the specification of an existing table; `IF NOT EXISTS` 
leaves an existing table unchanged.
+The two cannot be combined. `PARTITIONED BY` sets 
[`segmentGranularity`](#table-properties) and `CLUSTERED BY` sets
+`clusterKeys`, both of which a later `INSERT` or `REPLACE` inherits unless it 
states its own. `SEALED` sets
+[`sealed`](#table-properties), which requires every ingested column to be 
declared. It is written just before the
+column list because that is what it describes: the list is the table's whole 
schema, so the statement does not accept
+`SEALED` without one. (The `sealed` property itself can still be set on any 
table through `SET PROPERTIES` or the
+REST API.)
+
+Note that the table-level `CLUSTERED BY` is a sort order applied to each 
ingestion, which is a different thing from
+the `CLUSTERED BY` inside a [`__base` projection](#the-base-table), which 
defines how segments physically group rows.
+
+Column types are written as SQL types, such as `VARCHAR`, `BIGINT`, `DOUBLE`, 
or `VARCHAR ARRAY`. The `__time` column
+is written as `TIMESTAMP`. Types that have no SQL spelling, such as complex 
types, use `TYPE('...')` with the Druid
+native type string:
+
+```sql
+CREATE TABLE "druid"."visits" (
+  __time TIMESTAMP,
+  user_id VARCHAR,
+  pages_visited BIGINT,
+  sketch TYPE('COMPLEX<thetaSketch>')
+)
+PARTITIONED BY DAY
+CLUSTERED BY user_id
+```
+
+`ALTER TABLE` supports one change per statement, so that each statement is a 
single atomic catalog operation:
+
+```sql
+ALTER TABLE <table> ADD COLUMN <column> <type>
+ALTER TABLE <table> DROP COLUMN <column>
+ALTER TABLE <table> ALTER COLUMN <column> SET DATA TYPE <type>
+ALTER TABLE <table> ADD [IF NOT EXISTS] PROJECTION <name> AS ( <select> )
+ALTER TABLE <table> DROP PROJECTION [IF EXISTS] <name>
+ALTER TABLE <table> SET PROPERTIES ( <property> = <value> [, ...] )
+```
+
+`ADD COLUMN` fails if the column already exists, and `ALTER COLUMN` fails if 
it does not, so a misspelled column name
+is reported rather than quietly creating or replacing a column. Each statement 
is also checked against the rest of the
+table definition, not only the part it changes: adding a column, changing a 
type, or setting a property is rejected if
+the resulting table would be invalid, such as a segment granularity coarser 
than a projection the table declares.
+
+`DROP COLUMN` removes the column's declaration, not its data. A table's SQL 
schema is the declared columns followed by
+any other columns present in its segments, so a dropped column that still has 
data remains queryable: it loses its
+declared type and position and appears after the declared columns, like any 
column the catalog does not know about.
+Whether it can still be ingested then follows the usual rule: rejected if the 
table is sealed, accepted as an
+undeclared column otherwise. To remove a column from query results without 
rewriting data, use the catalog's
+`hiddenColumns` property instead, which currently has no SQL spelling and is 
set through the REST API.
+
+#### Projections
+
+A table may declare [projections](../../querying/projections.md), which are 
pre-aggregated views stored inside each
+segment. A projection is written as a `SELECT` over the table's own columns, 
with no `FROM` clause:
+
+```sql
+CREATE TABLE "druid"."visits" (
+  __time TIMESTAMP,
+  user_id VARCHAR,
+  user_agent VARCHAR,
+  pages_visited BIGINT,
+  PROJECTION daily_by_agent AS (
+    SELECT TIME_FLOOR(__time, 'P1D'), user_agent, SUM(pages_visited) AS 
total_pages
+    WHERE user_agent IS NOT NULL
+    GROUP BY 1, 2
+  )
+)
+PARTITIONED BY DAY
+```
+
+The body is planned exactly as the equivalent query would be, so a projection 
matches the queries it was written to
+serve. Every aggregate needs an alias, which becomes the name of the stored 
column. Time granularity is expressed
+with `TIME_FLOOR`, as it would be in a query.
+
+Because the body is planned like a query, it is planned under the statement's 
own query context, including any `SET`
+clauses. Only context parameters that affect planning can change the stored 
definition; parameters that only affect
+query execution have no effect, since the body is planned and stored rather 
than run. The context itself is not part
+of the definition: what the catalog stores is the projection the body planned 
to, so nothing from the statement's
+context is carried over to queries that later use it. Note also that a 
projection is matched to a query by its shape,
+so a definition planned under a context that changes that shape only matches 
queries run under the same context.
+
+A projection body accepts a select list, an optional `WHERE` and an optional 
`GROUP BY`. It cannot use `ORDER BY`,
+`LIMIT` or `HAVING`: a projection's ordering follows its grouping columns and 
is not something you choose. It also
+cannot use joins, subqueries, or expressions computed over aggregates, such as 
`SUM(x) / COUNT(x)`: store the two
+aggregates as separate columns and divide at query time. A projection can 
store the same
+[aggregation functions that are supported for rollup at ingestion 
time](../../multi-stage-query/concepts.md#rollup).
+
+Projections may also be added to and removed from an existing table:
+
+```sql
+ALTER TABLE "druid"."visits" ADD [IF NOT EXISTS] PROJECTION by_agent AS (
+  SELECT user_agent, SUM(pages_visited) AS total_pages GROUP BY user_agent
+)
+ALTER TABLE "druid"."visits" DROP PROJECTION [IF EXISTS] by_agent
+```
+
+Both take effect for subsequent ingestion. Segments already built keep 
whatever projections they were built with, so
+dropping a projection does not rewrite data.
+
+#### The base table
+
+The reserved projection name `__base` describes the table's own physical 
layout rather than an additional
+pre-aggregation. Defining it makes the table a 'clustered' table: rows of 
segments are stored grouped by the clustering
+columns.
+
+Its body lists the columns in the order segments store them, so it must name 
every declared column, in declared
+order. An item written as `<expr> AS <name>` makes that column computed at 
ingest time, from the columns it reads:
+
+```sql
+CREATE TABLE "druid"."events" SEALED (
+  tenant VARCHAR,
+  bucket BIGINT,
+  __time TIMESTAMP,
+  user_id BIGINT,
+  payload TYPE('COMPLEX<json>'),
+  PROJECTION __base AS (
+    SELECT tenant, ABS(user_id) % 128 AS bucket, __time, user_id, payload

Review Comment:
   I don't think `ABS(user_id) % 128` is valid Druid SQL. (Should use the `MOD` 
function.)



##########
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:
   Isn't it bad that time filters are retained as filters? Typically when a 
real SQL query is actually issued, `__time` filters will be moved to 
`intervals`, so they will not actually show up in the `filter`. So, does that 
mean they won't be able to match the stored projection? Or, is there something 
that causes this to be handled well?



-- 
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]

Reply via email to