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]