xiangfu0 commented on code in PR #19664:
URL: https://github.com/apache/pinot/pull/19664#discussion_r4201463248


##########
pinot-broker/src/main/java/org/apache/pinot/broker/broker/BasicAuthAccessControlFactory.java:
##########
@@ -133,6 +144,27 @@ public TableAuthorizationResult 
authorize(RequesterIdentity requesterIdentity, S
       return new TableAuthorizationResult(failedTables);
     }
 
+    @Override
+    public AuthorizationResult authorizeDeleteRows(RequesterIdentity 
requesterIdentity,

Review Comment:
   Added `authorizeDeleteRows` to `ZkBasicAuthAccessControlFactory` in 
d5065b2950. It matches the table by raw name, the way its query check already 
does for ZK users, and requires the explicit `DELETE` permission, with the same 
failure messages as here. The new `ZkBasicAuthAccessControlFactoryTest` covers 
a ZK broker user with `DELETE`, one with `READ` only, one with no permissions, 
missing or wrong credentials, a controller-only user, and the raw-name matching 
for typed and excluded tables.



##########
pinot-core/src/main/java/org/apache/pinot/core/query/executor/sql/SqlQueryExecutor.java:
##########
@@ -76,14 +86,45 @@ private static String getControllerBaseUrl(HelixManager 
helixManager) {
     return controllerBaseUrl;
   }
 
-  /// Execute DML Statement
+  /// Parses and executes a DML statement.
+  ///
+  /// A `DELETE` is refused by [#executeStatement] with a 
[QueryErrorCode#ACCESS_DENIED] error, since this method
+  /// cannot authorize the caller: it is executed once its table is resolved 
and the caller authorized to delete rows
+  /// from it, as the query endpoints of the broker and the controller do.
   ///
   /// @param sqlNodeAndOptions Parsed DML object
   /// @param headers extra headers map for minion task submission
   /// @return BrokerResponse is the DML executed response
   public BrokerResponse executeDMLStatement(SqlNodeAndOptions 
sqlNodeAndOptions,
       @Nullable Map<String, String> headers) {
-    DataManipulationStatement statement = 
DataManipulationStatementParser.parse(sqlNodeAndOptions);
+    DataManipulationStatement statement;
+    try {
+      statement = DataManipulationStatementParser.parse(sqlNodeAndOptions);
+    } catch (QueryException e) {
+      // e.g. a DELETE without a WHERE clause
+      return new BrokerResponseNative(e.getErrorCode(), e.getMessage());
+    }
+    return executeStatement(statement, headers);
+  }
+
+  /// Executes a parsed DML statement, e.g. from 
[DataManipulationStatementParser#parse].
+  ///
+  /// It does not authorize the caller. The table of a [DeleteStatement] must 
be resolved with
+  /// [DeleteStatement#resolveTableName] and the caller authorized to delete 
rows from it before it is executed, as the
+  /// query endpoints of the broker and the controller do: an unresolved 
`DELETE` is refused with a
+  /// [QueryErrorCode#ACCESS_DENIED] error.
+  ///
+  /// @param statement parsed statement
+  /// @param headers headers of the original request, e.g. for minion task 
submission
+  /// @return the response of the statement
+  public BrokerResponse executeStatement(DataManipulationStatement statement, 
@Nullable Map<String, String> headers) {
+    if (statement instanceof DeleteStatement) {
+      DeleteStatement deleteStatement = (DeleteStatement) statement;
+      if (!deleteStatement.isResolved()) {
+        return new BrokerResponseNative(QueryErrorCode.ACCESS_DENIED, 
UNAUTHORIZED_DELETE_MESSAGE);
+      }
+      return executeDelete(deleteStatement, headers);

Review Comment:
   Wrapped it in d5065b2950: the DELETE dispatch now lives in an `EXECUTOR` 
branch of the switch, a `QueryException` from `executeDelete` is returned under 
its own error code and anything else as `QUERY_EXECUTION`, and the Javadoc 
tells implementers they may throw `QueryException` with the code they want 
surfaced. On the controller the executor now runs right after authorization and 
only the serialization happens inside the `StreamingOutput`, so 
`executeSqlQueryCatching` maps anything that still escapes. 
`SqlQueryExecutorTest` has cases for an override throwing `QueryException` and 
one throwing a `RuntimeException`, and `PinotQueryResourceTest` has one with an 
executor that throws.



##########
pinot-common/src/main/java/org/apache/pinot/sql/parsers/dml/DeleteStatement.java:
##########
@@ -0,0 +1,307 @@
+/**
+ * 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.pinot.sql.parsers.dml;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.Set;
+import javax.annotation.Nullable;
+import org.apache.calcite.avatica.util.Casing;
+import org.apache.calcite.sql.SqlDelete;
+import org.apache.calcite.sql.SqlDialect;
+import org.apache.calcite.sql.SqlIdentifier;
+import org.apache.calcite.sql.SqlNode;
+import org.apache.calcite.sql.SqlOperator;
+import org.apache.calcite.sql.SqlSyntax;
+import org.apache.calcite.sql.fun.SqlStdOperatorTable;
+import org.apache.calcite.sql.parser.SqlParserPos;
+import org.apache.calcite.sql.util.SqlShuttle;
+import org.apache.calcite.sql.validate.SqlValidatorUtil;
+import org.apache.pinot.common.config.provider.TableCache;
+import org.apache.pinot.common.request.Expression;
+import org.apache.pinot.common.request.ExpressionType;
+import org.apache.pinot.common.request.Function;
+import org.apache.pinot.common.utils.DataSchema;
+import org.apache.pinot.common.utils.DatabaseUtils;
+import org.apache.pinot.spi.config.task.AdhocTaskConfig;
+import org.apache.pinot.spi.exception.QueryErrorCode;
+import org.apache.pinot.spi.exception.QueryException;
+import org.apache.pinot.spi.utils.CommonConstants;
+import org.apache.pinot.sql.parsers.CalciteSqlParser;
+import org.apache.pinot.sql.parsers.SqlNodeAndOptions;
+
+import static com.google.common.base.Preconditions.checkArgument;
+
+
+/// A SQL `DELETE FROM <table> WHERE <predicate>` statement.
+///
+/// Pinot parses `DELETE` but does not delete rows itself: the default SQL 
executor answers that it is not supported,
+/// and a deployment that can delete rows, e.g. by purging the matching rows 
from the segments with a minion task,
+/// executes the parsed statement by overriding 
`SqlQueryExecutor#executeDelete`. The generic [#execute()] and
+/// [#generateAdhocTaskConfig()] do not apply to it.
+///
+/// Before handing the statement to the executor, the broker and the 
controller resolve its table with
+/// [#resolveTableName] and authorize the caller to delete rows from that 
table.
+///
+/// Its options are the `SET` statements, the legacy `OPTION(...)` suffix and 
the request `queryOptions` of the
+/// statement (`SET` takes precedence). The `database` option only qualifies 
the table name, see [#resolveTableName].
+/// Every other option, including query options such as `timeoutMs`, reaches 
the executor as written through
+/// [#getOptions()], so that an option the executor relies on (e.g. a dry run) 
cannot be reclassified as a query option
+/// and silently dropped.
+///
+/// Instances are immutable and thread-safe.
+public class DeleteStatement implements DataManipulationStatement {
+  public static final String NOT_SUPPORTED_MESSAGE =
+      "DELETE is not supported by this Pinot cluster, it requires a SQL 
executor that implements row deletion";
+
+  /// Serializes the WHERE clause into SQL that the Pinot parser reads back 
into the same expression. Combined with
+  /// `quoteAllIdentifiers = false`, identifiers are only quoted (with double 
quotes) when they were quoted.
+  private static final SqlDialect PINOT_SQL_DIALECT = new 
SqlDialect(SqlDialect.EMPTY_CONTEXT
+      .withIdentifierQuoteString("\"")
+      .withLiteralQuoteString("'")
+      .withLiteralEscapedQuoteString("''")
+      .withUnquotedCasing(Casing.UNCHANGED)
+      .withQuotedCasing(Casing.UNCHANGED)
+      .withCaseSensitive(true));
+
+  /// Quotes the unquoted identifiers named like a SQL function without 
arguments, e.g. `user`, `pi` or `current_date`:
+  /// Calcite unparses them as upper-cased keywords, while Pinot reads them as 
columns, so they would name another
+  /// column. Mirrors the check of `SqlUtil.unparseSqlIdentifierSyntax`.
+  private static final SqlShuttle KEYWORD_IDENTIFIER_QUOTER = new SqlShuttle() 
{
+    @Override
+    public SqlNode visit(SqlIdentifier identifier) {
+      if (identifier.isSimple() && !identifier.getParserPosition().isQuoted()) 
{
+        SqlOperator operator =
+            
SqlValidatorUtil.lookupSqlFunctionByID(SqlStdOperatorTable.instance(), 
identifier, null);
+        if (operator != null && (operator.getSyntax() == SqlSyntax.FUNCTION_ID
+            || operator.getSyntax() == SqlSyntax.FUNCTION_ID_CONSTANT)) {
+          return new SqlIdentifier(identifier.names, null, 
SqlParserPos.QUOTED_ZERO,
+              List.of(SqlParserPos.QUOTED_ZERO));
+        }
+      }
+      return identifier;
+    }
+  };
+
+  /// Canonical names (lower case, without underscores) of the functions that 
read another table than the one the
+  /// statement deletes from, which the caller is not authorized to read.
+  private static final Set<String> CROSS_TABLE_FUNCTIONS = Set.of("lookup", 
"insubquery", "inpartitionedsubquery");

Review Comment:
   Done in d5065b2950: a `SCRIPT_FUNCTIONS` set with `groovy` is rejected by 
the same recursive walk with its own message, the `parse` and `getPredicate` 
Javadoc say so, and `DeleteStatementTest` covers a top-level `groovy(...)`, one 
under `AND`, and one nested inside `lower(...)`.



##########
pinot-broker/src/test/java/org/apache/pinot/broker/api/resources/PinotClientRequestTest.java:
##########
@@ -366,6 +395,358 @@ public void testProcessSqlQueryPostInvalidJson() {
     verify(_brokerMetrics, 
never()).addMeteredGlobalValue(BrokerMeter.UNCAUGHT_POST_EXCEPTIONS, 1L);
   }
 
+  @Test
+  public void testDmlOnGetQueryEndpointReturnsError()
+      throws Exception {
+    AsyncResponse asyncResponse = mock(AsyncResponse.class);
+    Request request = mock(Request.class);
+    when(request.getRequestURL()).thenReturn(new StringBuilder());
+    when(request.getHeaderNames()).thenReturn(List.of());
+
+    // A GET must not modify data, e.g. when a browser holding credentials 
follows a link
+    _pinotClientRequest.processSqlQueryGet("DELETE FROM myTable WHERE col1 = 
'a'", null, asyncResponse, request,
+        _httpHeaders);
+
+    ArgumentCaptor<Response> captor = ArgumentCaptor.forClass(Response.class);
+    verify(asyncResponse).resume(captor.capture());
+    
assertEquals(captor.getValue().getHeaders().get(PINOT_QUERY_ERROR_CODE_HEADER).get(0),
+        QueryErrorCode.SQL_PARSING.getId());
+    verify(_sqlQueryExecutor, never()).executeDMLStatement(any(), any());
+    verify(_sqlQueryExecutor, never()).executeStatement(any(), any());
+    verify(_requestHandler, never()).handleRequest(any(), any(), any(), any(), 
any());
+  }
+
+  @Test
+  public void testDeleteIsExecutedOnTheAuthorizedTable()
+      throws Exception {
+    RecordingAccessControl accessControl = new RecordingAccessControl();
+    when(_accessControlFactory.create()).thenReturn(accessControl);
+    // Table names are case-insensitive: the DELETE is authorized and executed 
on the table as it is defined
+    
when(_httpHeaders.getHeaderString(CommonConstants.DATABASE)).thenReturn("db1");
+    when(_tableCache.isIgnoreCase()).thenReturn(true);
+    
when(_tableCache.getActualTableName("db1.MYTABLE")).thenReturn("db1.myTable");
+    when(_sqlQueryExecutor.executeStatement(any(), any())).thenReturn(new 
BrokerResponseNative());
+
+    AsyncResponse asyncResponse = postSql("DELETE FROM MYTABLE WHERE col1 = 
'a'");
+
+    assertSucceeded(asyncResponse);
+    // The checks of a query on the table, then the right to delete rows
+    assertEquals(accessControl._checks, List.of("broker", "tables 
[db1.myTable]",
+        Actions.Table.QUERY + " TABLE db1.myTable", Actions.Table.DELETE_ROWS 
+ " TABLE db1.myTable",
+        "deleteRows db1.myTable"));
+    // The executor receives the table the caller is authorized for, and the 
headers to forward to the APIs it calls
+    DeleteStatement statement = executedDelete(Map.of("Authorization", "Basic 
abc"));
+    assertEquals(statement.getTableName(), "db1.myTable");
+    assertEquals(statement.getPredicate(), "col1 = 'a'");
+    verify(_sqlQueryExecutor, never()).executeDMLStatement(any(), any());
+    verify(_requestHandler, never()).handleRequest(any(), any(), any(), any(), 
any());
+  }
+
+  @Test
+  public void testDeleteOnTheMultiStageEndpointIsAuthorized()
+      throws Exception {
+    RecordingAccessControl accessControl = new RecordingAccessControl();
+    when(_accessControlFactory.create()).thenReturn(accessControl);
+    when(_sqlQueryExecutor.executeStatement(any(), any())).thenReturn(new 
BrokerResponseNative());
+
+    AsyncResponse asyncResponse = mock(AsyncResponse.class);
+    _pinotClientRequest.processSqlWithMultiStageQueryEnginePost(
+        JsonUtils.newObjectNode().put("sql", "DELETE FROM myTable WHERE col1 = 
'a'").toString(), asyncResponse, false,
+        0, mockRequest(), _httpHeaders);
+
+    assertSucceeded(asyncResponse);
+    assertEquals(accessControl._checks, List.of("broker", "tables [myTable]", 
Actions.Table.QUERY + " TABLE myTable",
+        Actions.Table.DELETE_ROWS + " TABLE myTable", "deleteRows myTable"));
+    assertEquals(executedDelete(Map.of("Authorization", "Basic 
abc")).getTableName(), "myTable");
+  }
+
+  @Test
+  public void testDeleteIsExecutedWithTheAllowAllAccessControl()
+      throws Exception {
+    when(_accessControlFactory.create()).thenReturn(new 
AllowAllAccessControlFactory().create());
+    when(_sqlQueryExecutor.executeStatement(any(), any())).thenReturn(new 
BrokerResponseNative());
+
+    assertSucceeded(postSql("DELETE FROM myTable WHERE col1 = 'a'"));
+    assertEquals(executedDelete(Map.of("Authorization", "Basic 
abc")).getTableName(), "myTable");
+  }
+
+  @Test
+  public void testDeleteWithoutTableCacheIsNotExecuted()
+      throws Exception {
+    // Broker applications using the legacy constructor do not bind a table 
cache.
+    FieldUtils.writeField(_pinotClientRequest, "_tableCache", null, true);
+    when(_accessControlFactory.create()).thenReturn(new 
AllowAllAccessControlFactory().create());
+
+    AsyncResponse asyncResponse = postSql("DELETE FROM myTable WHERE col1 = 
'a'");
+
+    ArgumentCaptor<Response> captor = ArgumentCaptor.forClass(Response.class);
+    verify(asyncResponse).resume(captor.capture());
+    
assertEquals(captor.getValue().getHeaders().get(PINOT_QUERY_ERROR_CODE_HEADER).get(0),
+        QueryErrorCode.QUERY_VALIDATION.getId());
+    verify(_sqlQueryExecutor, never()).executeStatement(any(), any());
+  }
+
+  @Test
+  public void testDeleteIsForbiddenByDefault()
+      throws Exception {
+    // An access control written before DELETE existed allows queries, but not 
deleting rows
+    when(_accessControlFactory.create()).thenReturn(new AccessControl() {
+      @Override
+      public TableAuthorizationResult authorize(RequesterIdentity 
requesterIdentity, Set<String> tables) {
+        return TableAuthorizationResult.success();
+      }
+    });
+
+    assertForbidden(postSql("DELETE FROM myTable WHERE col1 = 'a'"), "myTable",
+        "The access control of the broker does not allow deleting rows");
+  }
+
+  @DataProvider
+  public Object[][] deleteChecks() {
+    List<String> checks = List.of("broker", "tables [myTable]", 
Actions.Table.QUERY + " TABLE myTable",
+        Actions.Table.DELETE_ROWS + " TABLE myTable", "deleteRows myTable");
+    return new Object[][]{
+        {"tables", checks.subList(0, 2)},
+        {Actions.Table.QUERY, checks.subList(0, 3)},
+        {Actions.Table.DELETE_ROWS, checks.subList(0, 4)},
+        {"deleteRows", checks}
+    };
+  }
+
+  @Test(dataProvider = "deleteChecks")
+  public void testDeleteIsForbiddenByEachCheck(String deniedCheck, 
List<String> expectedChecks)
+      throws Exception {
+    RecordingAccessControl accessControl = new RecordingAccessControl(null, 
deniedCheck);
+    when(_accessControlFactory.create()).thenReturn(accessControl);
+
+    AsyncResponse asyncResponse = postSql("DELETE FROM myTable WHERE col1 = 
'a'");
+
+    assertForbidden(asyncResponse, "myTable", deniedCheck.equals("tables")
+        ? "Authorization Failed for tables: [myTable]"
+        : deniedCheck + " denied");
+    assertEquals(accessControl._checks, expectedChecks);
+  }
+
+  @Test
+  public void 
testDeleteFromATableWithATypeIsForbiddenByTheDeleteRowsActionOnItsRawName()

Review Comment:
   Added in d5065b2950: the positive `myTable_OFFLINE` test asserts the full 
check list (`QUERY` and `DeleteRows` on the typed name, `DeleteRows` on the raw 
name, then `authorizeDeleteRows`) and that the executor receives 
`myTable_OFFLINE`, plus a denial case that only denies `DeleteRows` on 
`myTable_OFFLINE`.



##########
pinot-segment-local/pom.xml:
##########
@@ -39,6 +39,11 @@
       <groupId>org.apache.pinot</groupId>
       <artifactId>pinot-common</artifactId>
     </dependency>
+    <dependency>

Review Comment:
   Dropped the pom hunk and the three test-harness changes from this PR in 
d5065b2950; the harness fixes are now their own PR, #19773, so they can land 
and be reverted on their own.



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