xiangfu0 commented on code in PR #19664:
URL: https://github.com/apache/pinot/pull/19664#discussion_r4191023928
##########
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:
The ZooKeeper flavour of this,
ZkBasicAuthAccessControlFactory.BasicAuthAccessControl, doesn't get a matching
override, so it inherits the deny-everything default and a broker configured
with it answers 403 on every DELETE no matter what the user was granted. That
bites because the ZK principal already carries the same permission list
(BasicAuthPrincipalUtils.extractBasicAuthPrincipals copies
UserConfig.getPermissions() in, and hasExplicitPermission comes from
BasicAuthPrincipal) and the controller's ZK factory honours it: a ZK user with
[READ, DELETE] on orders passes PinotQueryResource.authorizeDelete through
BaseBasicAuthAccessControl.hasAccess(raw, DELETE) but can never run the same
statement through the broker, so anyone managing users through ZK or the UI has
no in-tree way to turn on broker DELETE. The "written before DELETE existed"
rationale in the AccessControl Javadoc doesn't really fit a built-in factory
shipped right next to this one. I'd mirror this method there usin
g getPrincipalAuth(requesterIdentity), with the same hasTable(tableName) &&
hasTable(rawName) and hasExplicitPermission(AccessType.DELETE.name()) checks
and messages, plus a small test with a ZK BROKER user that has DELETE and one
that doesn't. If leaving it denied is deliberate, it's worth saying why in the
ZK factory's Javadoc and the release notes, since right now the only trace is
the "ZooKeeper basic auth denies it" line in the commit message.
##########
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:
This dispatch is the one path in `executeStatement` that lets an exception
escape: the MINION and HTTP branches right below catch `Exception` and return
it as a QUERY_EXECUTION error in the response, but an `executeDelete` override
that throws (a `QueryException` from predicate validation, which is how
`DeleteStatement.resolveTableName` reports errors, or an `IOException` from its
controller/minion call) goes straight through. On the broker that lands in
`processSqlQueryPost`'s catch-all, so the client gets an HTTP 500 with the bare
message, `UNCAUGHT_POST_EXCEPTIONS` ticks, and the full request JSON is logged
at ERROR; on the controller the executor runs inside the `StreamingOutput`
lambda in `PinotQueryResource.executeDelete`, after `executeSqlQueryCatching`
has already returned, so the same exception reaches
`WebApplicationExceptionMapper` and comes back as `{"code":500,...}` instead of
the `/sql` error payload that the identical `QueryException` thrown one frame
earlier by `re
solveTableName` or `authorizeDelete` produces. Since the Javadoc tells
implementers to validate the predicate and options but not how to report a
failure, I think the simplest fix is to wrap this call like its siblings and
say in the Javadoc that implementations may throw `QueryException` with the
error code they want surfaced:
```java
try {
return executeDelete(deleteStatement, headers);
} catch (QueryException e) {
return new BrokerResponseNative(e.getErrorCode(), e.getMessage());
} catch (Exception e) {
return new BrokerResponseNative(QueryErrorCode.QUERY_EXECUTION,
e.getMessage());
}
```
On the controller side it would also be worth calling `executeStatement`
eagerly next to `authorizeDelete` and only serializing the response in the
lambda, so anything that still escapes gets mapped by
`executeSqlQueryCatching`. A couple of `SqlQueryExecutorTest` cases with an
override that throws `QueryException` and one that throws `RuntimeException`,
asserting the error code in the returned response, would pin the contract down,
since all the executor stubs in the current tests return a response.
##########
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:
This is the only suffixed-table test that reaches the DELETE_ROWS checks,
and it only shows the raw-name check denying, so nothing in the broker tests
actually asserts that `DeleteRows` is consulted on `myTable_OFFLINE` itself
before falling through to `myTable`. If someone later aligns this with the
controller and checks DELETE_ROWS only on the raw name (or drops the suffixed
check in `authorizeDelete`), a fine-grained rule written against
`myTable_OFFLINE` would silently stop being consulted and every test here would
stay green, since all the full `_checks` assertions use an unsuffixed name
where the two checks collapse into one. The controller already has a positive
version of this in
`PinotQueryResourceTest.testDeleteFromATableWithATypeIsAuthorizedOnItsRawName`,
so it would be worth adding the mirror next to this one, something like:
```java
RecordingAccessControl accessControl = new RecordingAccessControl();
when(_accessControlFactory.create()).thenReturn(accessControl);
when(_sqlQueryExecutor.executeStatement(any(), any())).thenReturn(new
BrokerResponseNative());
assertSucceeded(postSql("DELETE FROM myTable_OFFLINE WHERE col1 = 'a'"));
assertEquals(accessControl._checks, List.of("broker", "tables
[myTable_OFFLINE]",
Actions.Table.QUERY + " TABLE myTable_OFFLINE",
Actions.Table.DELETE_ROWS + " TABLE myTable_OFFLINE",
Actions.Table.DELETE_ROWS + " TABLE myTable", "deleteRows
myTable_OFFLINE"));
assertEquals(executedDelete(Map.of("Authorization", "Basic
abc")).getTableName(), "myTable_OFFLINE");
```
A second denial case that denies `DeleteRows TABLE myTable_OFFLINE` by
description would also pin down that the suffixed check alone is enough to
refuse, mirroring the raw-name case you have here.
##########
pinot-segment-local/pom.xml:
##########
@@ -39,6 +39,11 @@
<groupId>org.apache.pinot</groupId>
<artifactId>pinot-common</artifactId>
</dependency>
+ <dependency>
Review Comment:
This org.jetbrains:annotations dependency doesn't look like it belongs to
the DELETE change: nothing under pinot-segment-local/src references it, master
builds this module without it (CI is green there), and every other module in
the repo only pulls it in at test scope, so this becomes the only
main-classpath declaration of that artifact. The "cannot access NotNull" error
from zstd-jni's Zstd class only shows up under the showWarnings/-Xlint style
compile that the precommit flow uses and reproduces on plain master, so it's a
pre-existing warning-flag quirk rather than something this branch introduced.
Bundling it here means whoever reviews a security-sensitive DML feature is also
signing off on an unexplained provided dep that dependency:analyze will flag as
unused-declared, and the next zstd-jni bump can't tell from this PR whether
it's safe to remove. I'd drop this hunk, and if you want the warning-flag
compile fix, send it on its own with an XML comment naming the zstd-jni @Not
Null reference and apply it consistently across the modules that call into
Zstd rather than just this one. Same thought for the
ControllerZkHelixUtilsTest, GroupByTrimmingTest and
BaseMultiClusterIntegrationTest changes: they're good fixes, but they're
independent CI stabilization, and if DELETE ever gets reverted to rework the
authorization model they'd disappear with it, so a small separate PR would let
them land and be reverted on their own and keep the release notes from
attributing them to SQL DELETE.
##########
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:
This set rejects lookUp and the IN_SUBQUERY variants at parse time because
an executor can't judge their safety, but `groovy(...)` sails through the same
check. For queries the broker refuses groovy in a filter by default
(`pinot.broker.disable.query.groovy=true`, with the per-table
`queryConfig.disableGroovy` override and the static analyzer when it is
enabled), yet that gate lives only in `BaseSingleStageBrokerRequestHandler`,
which a DELETE never enters: `BrokerRequestHandlerDelegate` rejects
`SqlDelete`, and neither `PinotClientRequest.executeDelete` nor
`PinotQueryResource.executeDelete` adds an equivalent. So `DELETE FROM t WHERE
groovy('{"returnType":"BOOLEAN","isSingleValue":true}', '<script>', col)`
parses, round-trips, passes `authorizeDelete` and lands in `executeDelete` as a
predicate string, and an executor that evaluates it on a minion or server has
none of the broker config, the table override, or the analyzer
(`configureGroovySecurity` is only called by the broker
and controller starters), so any principal with DELETE on the table gets an
unrestricted script run there. The stock executor refuses every DELETE, so
nothing shipped here executes it, but this is exactly the extension point the
PR adds and the `executeDelete`/`getPredicate` Javadoc only mentions columns
and aggregations. Since you already reject functions here for policy reasons,
I'd add groovy to the same list (or a sibling set with its own message) and a
`DeleteStatementTest` case next to the lookUp/IN_SUBQUERY ones, and say so in
the executor Javadoc so implementers know what has and hasn't been validated:
```java
private static final Set<String> CROSS_TABLE_FUNCTIONS = Set.of("lookup",
"insubquery", "inpartitionedsubquery");
/// Rejected because the broker's groovy policy (DISABLE_GROOVY, per-table
override, static analyzer) is not applied to DELETE.
private static final Set<String> SCRIPT_FUNCTIONS = Set.of("groovy");
```
##########
pinot-controller/src/main/java/org/apache/pinot/controller/api/resources/PinotQueryResource.java:
##########
@@ -424,6 +439,43 @@ private StreamingOutput executeSqlQuery(@Context
HttpHeaders httpHeaders, String
}
}
+ /// Executes a `DELETE` once the caller is authorized to delete rows from
its table, see [#authorizeDelete]. The
+ /// table is resolved first, with the database of the request and in the
case it is defined with, and the executor
+ /// deletes rows from that exact table.
+ private StreamingOutput executeDelete(SqlNodeAndOptions sqlNodeAndOptions,
HttpHeaders httpHeaders) {
+ DeleteStatement statement = ((DeleteStatement)
DataManipulationStatementParser.parse(sqlNodeAndOptions))
Review Comment:
Minor, but this resolves the table before any access check runs, and since
handlePostSql is @ManualAuthorization nothing has authenticated the caller yet.
If the name happens to be a logical table, resolveTableName throws
QUERY_VALIDATION "DELETE does not support logical tables: <name>" and
executeSqlQueryCatching returns that as a 200, whereas any other name falls
through to hasAccess and, under BasicAuth, gets a 401 from the missing
principal. So an unauthenticated client can tell logical-table names apart from
everything else on /sql, which GET /logicalTables would not let it do, and it
is also out of step with the MSE path here (hasAccess before touching the
cache) and the broker's DELETE path (first-step authorize before the lookup,
which the PinotClientRequestTest comment calls out for exactly this reason). A
caller-level check at the top of executeDelete, the same one
getMultiStageQueryResponse uses, would close it with BasicAuth throwing 401
before the TableCache is consul
ted, and then the per-table authorizeDelete can stay as it is:
```java
if (!_accessControlFactory.create().hasAccess(AccessType.READ, httpHeaders,
SQL_ENDPOINT)) {
throw QueryErrorCode.ACCESS_DENIED.asException("Permission denied to
delete rows");
}
```
It would be worth a test next to testUnauthenticatedDeleteIsUnauthorized
where the TableCache says the name is a logical table and the no-table
hasAccess throws NotAuthorizedException, asserting the 401 wins and
getActualLogicalTableName is never called.
##########
pinot-common/src/main/java/org/apache/pinot/sql/parsers/CalciteSqlParser.java:
##########
@@ -678,6 +682,11 @@ public static Expression compileToExpression(String
expression) {
return toExpression(sqlNode);
}
+ /// Compiles an expression that is already parsed, e.g. the condition of a
parsed statement, into [Expression].
+ public static Expression compileToExpression(SqlNode sqlNode) {
Review Comment:
This new overload quietly has different semantics from
`compileToExpression(String)` right above it: the String version runs
`PostgreSqlCastRewriter.rewrite` before `toExpression`, this one doesn't. So
for the same SQL, `compileToExpression("intCol::double > 1")` throws while
handing the parsed node here returns a `cast(intCol,'DOUBLE') > 1`, and
`'\x01'::bytea` comes out as `X'01'` from one and as a raw
`cast('\x01','BYTEA')` call from the other. Nothing in the PR trips over it
because `DeleteStatement` only ever gets a condition that already went through
`compileToSqlNodeAndOptions` (and `toPinotSql` fails closed on a mismatch
anyway), but the rewriter is package-private, so an outside caller who parses a
node themselves has no way to apply the same normalization, and the Javadoc
doesn't say the node is expected to come from `compileToSqlNodeAndOptions`. The
cheapest fix is to just run the rewriter here too, which is a no-op on a tree
that's already been rewritten since no `::`
casts survive the first pass and `rewriteCast` returns ordinary CAST calls
unchanged:
```java
public static Expression compileToExpression(SqlNode sqlNode) {
return toExpression(PostgreSqlCastRewriter.rewrite(sqlNode));
}
```
If you'd rather keep it as-is, I'd at least state the precondition in the
Javadoc so the two overloads don't look interchangeable when they aren't.
##########
pinot-broker/src/main/java/org/apache/pinot/broker/api/resources/PinotClientRequest.java:
##########
@@ -713,6 +738,112 @@ private BrokerResponse executeSqlQuery(ObjectNode
sqlRequestJson, HttpRequesterI
}
}
+ /// Executes a `DELETE` once the caller is authorized to delete rows from
its table.
+ ///
+ /// The first-step access control runs first, as for queries. The table is
then resolved with the database of the
+ /// request and in the case it is defined with, the caller is authorized to
delete rows from it (see
+ /// [#authorizeDelete]), and the executor deletes rows from that exact table.
+ private BrokerResponse executeDelete(SqlNodeAndOptions sqlNodeAndOptions,
Map<String, String> headers,
+ HttpRequesterIdentity requesterIdentity, @Nullable HttpHeaders
httpHeaders) {
+ AccessControl accessControl = _accessControlFactory.create();
+ // The first-step access control runs before the table is looked up, as
for queries
+ AuthorizationResult authorizationResult =
accessControl.authorize(requesterIdentity);
+ if (!authorizationResult.hasAccess()) {
+ throw deleteAccessDenied(null, authorizationResult);
+ }
+ if (_tableCache == null) {
+ return new BrokerResponseNative(QueryErrorCode.QUERY_VALIDATION,
+ "DELETE is not supported by this broker: no table cache was
configured");
+ }
+ DeleteStatement statement;
+ try {
+ String databaseHeader = httpHeaders != null ?
httpHeaders.getHeaderString(CommonConstants.DATABASE) : null;
+ statement = ((DeleteStatement)
DataManipulationStatementParser.parse(sqlNodeAndOptions))
+ .resolveTableName(databaseHeader, _tableCache);
+ } catch (QueryException e) {
Review Comment:
This catch is the only thing keeping a malformed DELETE out of the generic
500 path in processSqlQueryPost, but no broker test actually reaches it: every
DELETE posted in PinotClientRequestTest has a WHERE clause, and the two
logical-table cases are denied by the first-step authorize before parse or
resolveTableName ever runs. So if someone later narrows this to
DatabaseConflictException, or parse/resolve starts throwing something that
isn't a QueryException, `DELETE FROM myTable` from an authorized caller would
come back as HTTP 500 with UNCAUGHT_POST_EXCEPTIONS bumped instead of a 200
body carrying SQL_PARSING, and the suite would stay green. The parse and
resolve failures themselves are covered in DeleteStatementTest, so this is only
about pinning the endpoint's mapping. The existing postSql helper plus the
AllowAll stub make it a few lines each; a missing-WHERE case, a logical-table
case with the AllowAll factory, and a `database` header that conflicts with
`DELETE FROM db1.my
Table ...` would cover the three cases the comment on this catch lists:
```java
when(_accessControlFactory.create()).thenReturn(new
AllowAllAccessControlFactory().create());
AsyncResponse asyncResponse = postSql("DELETE FROM myTable");
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()).executeStatement(any(), any());
verify(_brokerMetrics,
never()).addMeteredGlobalValue(BrokerMeter.UNCAUGHT_POST_EXCEPTIONS, 1L);
```
The missing-WHERE and database-conflict cases are just as cheap to mirror in
PinotQueryResourceTest.
##########
pinot-broker/src/test/java/org/apache/pinot/broker/broker/BasicAuthAccessControlTest.java:
##########
@@ -199,4 +214,44 @@ public void testNormalizeToken() {
Assert.assertTrue(_accessControl.authorize(identity, request).hasAccess());
Assert.assertTrue(_accessControl.authorize(identity,
_tableNames).hasAccess());
}
+
+ @Test
+ public void testDeleteRowsRequiresTheDeletePermission() {
Review Comment:
The new class Javadoc on BasicAuthAccessControlFactory promises that a
principal configured with `permissions=*` cannot delete rows, but nothing in
this test pins that. It currently holds only because
BasicAuthPrincipalUtils.extractSet collapses a bare `*` to an empty set and
hasExplicitPermission is a plain contains(), and neither of those is exercised
with `*` anywhere in the suite (the fixture only has no-permissions, `read`,
and `read,delete` principals). Since `*` is the usual way people configure an
admin, a later change that turns it into a real wildcard to match the
controller's semantics would silently grant DELETE to every such principal
while this test keeps passing. Could you add a `*` principal to the fixture and
assert it here, something like:
```java
config.put("principals.star.permissions", "*");
...
assertDeleteRows(TOKEN_STAR, "lessImportantStuff", false, "Principal: star
is not granted the DELETE permission");
```
A `read,DELETE` principal asserted as allowed would also be cheap and would
pin the case-insensitive match.
##########
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");
+
+ private final String _tableName;
+ private final String _predicate;
+ @Nullable
+ private final String _database;
+ private final Map<String, String> _options;
+ private final boolean _resolved;
+
+ private DeleteStatement(String tableName, String predicate, @Nullable String
database, Map<String, String> options) {
+ this(tableName, predicate, database, options, false);
+ }
+
+ private DeleteStatement(String tableName, String predicate, @Nullable String
database, Map<String, String> options,
+ boolean resolved) {
+ _tableName = tableName;
+ _predicate = predicate;
+ _database = database;
+ _options = Collections.unmodifiableMap(new HashMap<>(options));
+ _resolved = resolved;
+ }
+
+ /// Parses a `DELETE` statement.
+ ///
+ /// @throws IllegalArgumentException if the statement is not a supported
`DELETE`: it must have a WHERE clause, no
+ /// table alias, a WHERE clause that Pinot
can parse as an expression and that does
+ /// not read another table (e.g. with
`lookUp` or `IN_SUBQUERY`), and set the
+ /// database with the `database` option only
+ public static DeleteStatement parse(SqlNodeAndOptions sqlNodeAndOptions) {
+ SqlNode sqlNode = sqlNodeAndOptions.getSqlNode();
+ checkArgument(sqlNode instanceof SqlDelete, "Not a DELETE statement: %s",
sqlNode.getKind());
+ SqlDelete sqlDelete = (SqlDelete) sqlNode;
+ String tableName = getTableName(sqlDelete.getTargetTable());
+ checkArgument(sqlDelete.getAlias() == null,
+ "DELETE does not support a table alias, reference the columns
directly");
+ SqlNode condition = sqlDelete.getCondition();
+ checkArgument(condition != null,
+ "DELETE requires a WHERE clause; delete the table segments to remove
all of its rows");
+ String predicate = toPinotSql(condition);
+
+ String database = null;
+ Map<String, String> options = new HashMap<>();
+ for (Map.Entry<String, String> option :
sqlNodeAndOptions.getOptions().entrySet()) {
+ String key = option.getKey();
+ if (key.equals(CommonConstants.DATABASE)) {
+ database = option.getValue();
+ } else if (key.equalsIgnoreCase(CommonConstants.DATABASE)) {
+ // Queries ignore it, as they only read the `database` option: fail
rather than delete from another table
+ throw new IllegalArgumentException(
+ "Unsupported option: " + key + ", set the database with the '" +
CommonConstants.DATABASE + "' option");
+ } else {
+ options.put(key, option.getValue());
+ }
+ }
+ return new DeleteStatement(tableName, predicate, database, options);
+ }
+
+ private static String getTableName(SqlNode targetTable) {
+ // Table hints and EXTEND clauses parse into other node types
+ checkArgument(targetTable instanceof SqlIdentifier, "DELETE only supports
a plain table name, got: %s",
+ targetTable);
+ // A quoted name part may contain a dot, which splits it like the table
name of a query
+ String tableName = String.join(".", ((SqlIdentifier) targetTable).names);
+ // Empty parts, e.g. in "db."."t", are rejected rather than dropped, so
that the table name is the one authorized
+ String[] parts = tableName.split("\\.", -1);
+ checkArgument(parts.length <= 2 &&
Arrays.stream(parts).noneMatch(String::isEmpty),
+ "Invalid table name: %s, expected [database.]table", tableName);
+ return tableName;
+ }
+
+ /// Serializes a WHERE clause back into SQL, and verifies that Pinot parses
it into the same expression.
+ private static String toPinotSql(SqlNode condition) {
+ Expression expression;
+ Expression serializedExpression;
+ String predicate;
+ try {
+ expression = CalciteSqlParser.compileToExpression(condition);
+ predicate = condition.accept(KEYWORD_IDENTIFIER_QUOTER)
+ .toSqlString(config -> config.withDialect(PINOT_SQL_DIALECT)
+ .withQuoteAllIdentifiers(false)
+ .withIndentation(0))
+ .getSql();
+ serializedExpression = CalciteSqlParser.compileToExpression(predicate);
+ } catch (Exception e) {
+ throw new IllegalArgumentException("Unsupported WHERE clause in DELETE:
" + condition, e);
+ }
+ // Fail rather than hand over a predicate that selects other rows than the
statement
+ checkArgument(serializedExpression.equals(expression),
+ "Unsupported WHERE clause in DELETE, it cannot be serialized back into
the same expression: %s", condition);
+ checkNoCrossTableFunction(expression);
+ return predicate;
+ }
+
+ /// Rejects the functions that read another table: the caller is only
authorized for the table it deletes from.
+ private static void checkNoCrossTableFunction(Expression expression) {
+ if (expression.getType() != ExpressionType.FUNCTION) {
+ return;
+ }
+ Function function = expression.getFunctionCall();
+ String functionName = function.getOperator().replace("_",
"").toLowerCase(Locale.ROOT);
+ checkArgument(!CROSS_TABLE_FUNCTIONS.contains(functionName),
+ "Unsupported WHERE clause in DELETE, %s reads another table",
function.getOperator());
+ if (function.getOperands() != null) {
+ for (Expression operand : function.getOperands()) {
+ checkNoCrossTableFunction(operand);
+ }
+ }
+ }
+
+ /// Returns the statement with its table name resolved: qualified with the
database of the request, and in the case
+ /// the table is defined with (table names are case-insensitive by default).
The broker and the controller authorize
+ /// the caller for the resolved table name, and hand the resolved statement
to the executor.
+ ///
+ /// The database of the request is the `database` request header, else the
`database` option of the statement, as
+ /// for a multi-stage query (see
`DatabaseUtils#extractDatabaseFromQueryRequest`). They must match when both are
set,
+ /// and the database of a `database.table` name must match them.
+ ///
+ /// @param databaseHeader value of the `database` request header, if any
+ /// @param tableCache tables of the cluster, to resolve the case of the
table name. A table name it does not know
+ /// keeps the case of the statement.
+ /// @throws QueryException with [QueryErrorCode#QUERY_VALIDATION] if the
`database` header, the `database` option
+ /// and the database of a `database.table` name do
not match (a
+ /// `DatabaseConflictException`), or if the table is
a logical table, which `DELETE` does
+ /// not support
+ public DeleteStatement resolveTableName(@Nullable String databaseHeader,
TableCache tableCache)
+ throws QueryException {
+ String database =
DatabaseUtils.extractDatabaseFromOptionAndHeader(_database, databaseHeader);
+ String tableName;
+ try {
+ tableName = DatabaseUtils.translateTableName(_tableName, database,
tableCache.isIgnoreCase());
+ } catch (IllegalArgumentException e) {
+ throw QueryErrorCode.QUERY_VALIDATION.asException("Invalid table name in
DELETE: " + e.getMessage(), e);
+ }
+ String actualTableName = tableCache.getActualTableName(tableName);
+ if (actualTableName == null &&
tableCache.getActualLogicalTableName(tableName) != null) {
+ // Deleting from a logical table would delete from physical tables the
caller is not authorized for
+ throw QueryErrorCode.QUERY_VALIDATION.asException("DELETE does not
support logical tables: " + tableName);
+ }
+ return new DeleteStatement(actualTableName != null ? actualTableName :
tableName, _predicate, null, _options, true);
Review Comment:
On a TableCache miss this marks the statement resolved with the name exactly
as the caller typed it, and both PinotClientRequest.executeDelete and
PinotQueryResource.executeDelete then authorize that string and hand it to the
executor with isResolved()=true, so there is no point where an unknown table
fails. The query path on the same miss authorizes and then returns
TABLE_DOES_NOT_EXIST before dispatch, and the asymmetry matters because
BasicAuthPrincipal.hasTable compares include/exclude sets case-sensitively
while the default cache is case-insensitive: with excludeTables=Secret, `DELETE
FROM SECRET` during a cache-miss window (say, before the ZK child-change
callback lands after Secret is created) is authorized as "SECRET", whereas the
warm-cache resolution to "Secret" is denied. Nothing in-tree deletes today
since the default executor returns not-supported, and the Javadoc on
getTableName tells executors not to re-resolve, but right now that note is the
only thing standing bet
ween an executor that resolves case-insensitively and a cross-table delete,
and it also makes the authorizeDeleteRows promise of a name "in the case the
table is defined with" false for a miss. I'd keep the authorization on the
case-resolved name as you have it (so existence isn't leaked), but record
whether the cache actually knew the table and fail closed afterwards with
QueryErrorCode.TABLE_DOES_NOT_EXIST, matching the query path, and adjust the
"db1.OtherTable" case in DeleteStatementTest accordingly.
##########
pinot-controller/src/main/java/org/apache/pinot/controller/api/resources/PinotQueryResource.java:
##########
@@ -397,6 +405,10 @@ private StreamingOutput executeSqlQuery(@Context
HttpHeaders httpHeaders, String
throw QueryErrorCode.QUERY_VALIDATION.asException(
"DDL statements are not supported on /sql; use POST /sql/ddl
instead.");
}
+ if (isGet && sqlNodeAndOptions.getSqlNode() instanceof SqlDelete) {
Review Comment:
Small consistency thing on this GET guard: it only keys on `SqlDelete`, so
`GET /sql?sql=INSERT INTO t FROM FILE ...` still falls through to the `case
DML` branch below and submits the ingestion task, even though the "browser
following a link that mutates data" rationale in the `handleGetSql` comment
applies just as much to INSERT. The broker's `GET /query/sql` already rejects
every non-DQL statement via `onlyDql && sqlType != PinotSqlType.DQL` (with
SQL_PARSING), so right now the same GET DELETE gets QUERY_VALIDATION from the
controller and SQL_PARSING from the broker, and GET INSERT is rejected on one
role but executed on the other. I know INSERT-on-GET predates this PR, but
since this is where the GET gating gets introduced, it seems cheap to make the
controller match the broker:
```java
if (isGet && sqlType != PinotSqlType.DQL) {
throw QueryErrorCode.SQL_PARSING.asException(
"Unsupported SQL type - " + sqlType + ", GET /sql only supports DQL;
use POST /sql instead.");
}
```
That would mean adjusting
`testDeleteOnGetQueryEndpointReturnsValidationError` for the chosen code and
adding a GET INSERT case that verifies `executeDMLStatement` is never called.
If you'd rather keep this DELETE-only to avoid touching the INSERT path here, a
short note in the `handleGetSql` comment saying so would make the asymmetry
read as deliberate.
##########
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");
+
+ private final String _tableName;
+ private final String _predicate;
+ @Nullable
+ private final String _database;
+ private final Map<String, String> _options;
+ private final boolean _resolved;
+
+ private DeleteStatement(String tableName, String predicate, @Nullable String
database, Map<String, String> options) {
+ this(tableName, predicate, database, options, false);
+ }
+
+ private DeleteStatement(String tableName, String predicate, @Nullable String
database, Map<String, String> options,
+ boolean resolved) {
+ _tableName = tableName;
+ _predicate = predicate;
+ _database = database;
+ _options = Collections.unmodifiableMap(new HashMap<>(options));
+ _resolved = resolved;
+ }
+
+ /// Parses a `DELETE` statement.
+ ///
+ /// @throws IllegalArgumentException if the statement is not a supported
`DELETE`: it must have a WHERE clause, no
+ /// table alias, a WHERE clause that Pinot
can parse as an expression and that does
+ /// not read another table (e.g. with
`lookUp` or `IN_SUBQUERY`), and set the
+ /// database with the `database` option only
+ public static DeleteStatement parse(SqlNodeAndOptions sqlNodeAndOptions) {
+ SqlNode sqlNode = sqlNodeAndOptions.getSqlNode();
+ checkArgument(sqlNode instanceof SqlDelete, "Not a DELETE statement: %s",
sqlNode.getKind());
+ SqlDelete sqlDelete = (SqlDelete) sqlNode;
+ String tableName = getTableName(sqlDelete.getTargetTable());
+ checkArgument(sqlDelete.getAlias() == null,
+ "DELETE does not support a table alias, reference the columns
directly");
+ SqlNode condition = sqlDelete.getCondition();
+ checkArgument(condition != null,
+ "DELETE requires a WHERE clause; delete the table segments to remove
all of its rows");
+ String predicate = toPinotSql(condition);
+
+ String database = null;
+ Map<String, String> options = new HashMap<>();
+ for (Map.Entry<String, String> option :
sqlNodeAndOptions.getOptions().entrySet()) {
+ String key = option.getKey();
+ if (key.equals(CommonConstants.DATABASE)) {
+ database = option.getValue();
+ } else if (key.equalsIgnoreCase(CommonConstants.DATABASE)) {
+ // Queries ignore it, as they only read the `database` option: fail
rather than delete from another table
+ throw new IllegalArgumentException(
+ "Unsupported option: " + key + ", set the database with the '" +
CommonConstants.DATABASE + "' option");
+ } else {
+ options.put(key, option.getValue());
+ }
+ }
+ return new DeleteStatement(tableName, predicate, database, options);
+ }
+
+ private static String getTableName(SqlNode targetTable) {
+ // Table hints and EXTEND clauses parse into other node types
+ checkArgument(targetTable instanceof SqlIdentifier, "DELETE only supports
a plain table name, got: %s",
+ targetTable);
+ // A quoted name part may contain a dot, which splits it like the table
name of a query
+ String tableName = String.join(".", ((SqlIdentifier) targetTable).names);
+ // Empty parts, e.g. in "db."."t", are rejected rather than dropped, so
that the table name is the one authorized
+ String[] parts = tableName.split("\\.", -1);
+ checkArgument(parts.length <= 2 &&
Arrays.stream(parts).noneMatch(String::isEmpty),
+ "Invalid table name: %s, expected [database.]table", tableName);
+ return tableName;
+ }
+
+ /// Serializes a WHERE clause back into SQL, and verifies that Pinot parses
it into the same expression.
+ private static String toPinotSql(SqlNode condition) {
+ Expression expression;
+ Expression serializedExpression;
+ String predicate;
+ try {
+ expression = CalciteSqlParser.compileToExpression(condition);
+ predicate = condition.accept(KEYWORD_IDENTIFIER_QUOTER)
+ .toSqlString(config -> config.withDialect(PINOT_SQL_DIALECT)
+ .withQuoteAllIdentifiers(false)
+ .withIndentation(0))
+ .getSql();
+ serializedExpression = CalciteSqlParser.compileToExpression(predicate);
+ } catch (Exception e) {
+ throw new IllegalArgumentException("Unsupported WHERE clause in DELETE:
" + condition, e);
Review Comment:
This catch block unparses the node a second time via `"... " + condition`
(SqlNode.toString), so if the unparse in the try block is what failed, the same
exception fires again while building the message and the
IllegalArgumentException is never constructed. That happens today with `DELETE
FROM myTable WHERE ts AT TIME ZONE 'pst' > 123`: SqlAtTimeZone is a
SqlSpecialOperator with no unparse override, so Calcite throws
UnsupportedOperationException("class org.apache.calcite.sql.SqlSyntax$7:
SPECIAL"), which escapes parse(), skips the SQL_PARSING mapping in
DataManipulationStatementParser, and surfaces as an HTTP 500 with that opaque
body plus an ERROR log and UNCAUGHT_POST_EXCEPTIONS bump on the broker (and an
INTERNAL error on the controller), while every other bad WHERE clause gets a
clean SQL_PARSING error. Since the cause already carries the detail, the
simplest fix is to not re-render the node here, or to format it leniently the
way the checkArgument a few lines down already do
es:
```java
throw new IllegalArgumentException(
Strings.lenientFormat("Unsupported WHERE clause in DELETE: %s",
condition), e);
```
It would be worth adding the AT TIME ZONE case to DeleteStatementTest
asserting the IllegalArgumentException so this path stays SQL_PARSING on both
roles.
##########
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,
+ @Nullable HttpHeaders httpHeaders, String tableName) {
+ Optional<BasicAuthPrincipal> principalOpt =
getPrincipalOpt(requesterIdentity);
+ if (principalOpt.isEmpty()) {
+ return new BasicAuthorizationResultImpl(false, "Missing or invalid
credentials");
+ }
+ BasicAuthPrincipal principal = principalOpt.get();
+ // The table checks match the name as given, while `excludeTables` lists
raw names: check both, so that a name
+ // with a type suffix cannot bypass an excluded table
+ if (!principal.hasTable(tableName) ||
!principal.hasTable(TableNameBuilder.extractRawTableName(tableName))) {
Review Comment:
The second `hasTable` call does more than the comment says:
`BasicAuthPrincipal.hasTable` is `isTableIncluded && isTableNotExcluded`, so
calling it on the raw name re-applies the `tables` inclusion list to the raw
name, not just `excludeTables`. That breaks a principal configured as
`tables=events_REALTIME, permissions=read,delete`: `SELECT * FROM
events_REALTIME` is authorized because `authorize` only matches the name as
given, and `DELETE FROM events_REALTIME WHERE ...` passes `authorize(identity,
Set.of("events_REALTIME"))` and the fine-grained checks, then gets denied here
with "does not have access to table: events_REALTIME" because `_tables` does
not contain `events`. The only workaround is to also list `events` in `tables`,
which widens that principal's query access to the hybrid name. I'd split the
two concerns, something like an `isTableExcluded(String)` on
`BasicAuthPrincipal` so this reads
```java
if (!principal.hasTable(tableName)
||
!principal.isTableExcluded(TableNameBuilder.extractRawTableName(tableName))) {
```
(with the negation the right way around), and add a fixture principal with
`tables=lessImportantStuff_OFFLINE, permissions=read,delete` to
`BasicAuthAccessControlTest` asserting the typed name is allowed and the raw
name is denied, mirroring the existing `deleter` case in
`testDeleteRowsChecksTheRawTableName`.
##########
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");
+
+ private final String _tableName;
+ private final String _predicate;
+ @Nullable
+ private final String _database;
+ private final Map<String, String> _options;
+ private final boolean _resolved;
+
+ private DeleteStatement(String tableName, String predicate, @Nullable String
database, Map<String, String> options) {
+ this(tableName, predicate, database, options, false);
+ }
+
+ private DeleteStatement(String tableName, String predicate, @Nullable String
database, Map<String, String> options,
+ boolean resolved) {
+ _tableName = tableName;
+ _predicate = predicate;
+ _database = database;
+ _options = Collections.unmodifiableMap(new HashMap<>(options));
+ _resolved = resolved;
+ }
+
+ /// Parses a `DELETE` statement.
+ ///
+ /// @throws IllegalArgumentException if the statement is not a supported
`DELETE`: it must have a WHERE clause, no
+ /// table alias, a WHERE clause that Pinot
can parse as an expression and that does
+ /// not read another table (e.g. with
`lookUp` or `IN_SUBQUERY`), and set the
+ /// database with the `database` option only
+ public static DeleteStatement parse(SqlNodeAndOptions sqlNodeAndOptions) {
+ SqlNode sqlNode = sqlNodeAndOptions.getSqlNode();
+ checkArgument(sqlNode instanceof SqlDelete, "Not a DELETE statement: %s",
sqlNode.getKind());
+ SqlDelete sqlDelete = (SqlDelete) sqlNode;
+ String tableName = getTableName(sqlDelete.getTargetTable());
+ checkArgument(sqlDelete.getAlias() == null,
+ "DELETE does not support a table alias, reference the columns
directly");
+ SqlNode condition = sqlDelete.getCondition();
+ checkArgument(condition != null,
+ "DELETE requires a WHERE clause; delete the table segments to remove
all of its rows");
+ String predicate = toPinotSql(condition);
+
+ String database = null;
+ Map<String, String> options = new HashMap<>();
+ for (Map.Entry<String, String> option :
sqlNodeAndOptions.getOptions().entrySet()) {
+ String key = option.getKey();
+ if (key.equals(CommonConstants.DATABASE)) {
+ database = option.getValue();
+ } else if (key.equalsIgnoreCase(CommonConstants.DATABASE)) {
+ // Queries ignore it, as they only read the `database` option: fail
rather than delete from another table
+ throw new IllegalArgumentException(
+ "Unsupported option: " + key + ", set the database with the '" +
CommonConstants.DATABASE + "' option");
+ } else {
+ options.put(key, option.getValue());
+ }
+ }
+ return new DeleteStatement(tableName, predicate, database, options);
+ }
+
+ private static String getTableName(SqlNode targetTable) {
+ // Table hints and EXTEND clauses parse into other node types
+ checkArgument(targetTable instanceof SqlIdentifier, "DELETE only supports
a plain table name, got: %s",
+ targetTable);
+ // A quoted name part may contain a dot, which splits it like the table
name of a query
+ String tableName = String.join(".", ((SqlIdentifier) targetTable).names);
+ // Empty parts, e.g. in "db."."t", are rejected rather than dropped, so
that the table name is the one authorized
+ String[] parts = tableName.split("\\.", -1);
+ checkArgument(parts.length <= 2 &&
Arrays.stream(parts).noneMatch(String::isEmpty),
+ "Invalid table name: %s, expected [database.]table", tableName);
+ return tableName;
+ }
+
+ /// Serializes a WHERE clause back into SQL, and verifies that Pinot parses
it into the same expression.
+ private static String toPinotSql(SqlNode condition) {
+ Expression expression;
+ Expression serializedExpression;
+ String predicate;
+ try {
+ expression = CalciteSqlParser.compileToExpression(condition);
+ predicate = condition.accept(KEYWORD_IDENTIFIER_QUOTER)
+ .toSqlString(config -> config.withDialect(PINOT_SQL_DIALECT)
+ .withQuoteAllIdentifiers(false)
+ .withIndentation(0))
+ .getSql();
+ serializedExpression = CalciteSqlParser.compileToExpression(predicate);
+ } catch (Exception e) {
+ throw new IllegalArgumentException("Unsupported WHERE clause in DELETE:
" + condition, e);
+ }
+ // Fail rather than hand over a predicate that selects other rows than the
statement
+ checkArgument(serializedExpression.equals(expression),
+ "Unsupported WHERE clause in DELETE, it cannot be serialized back into
the same expression: %s", condition);
+ checkNoCrossTableFunction(expression);
+ return predicate;
+ }
+
+ /// Rejects the functions that read another table: the caller is only
authorized for the table it deletes from.
+ private static void checkNoCrossTableFunction(Expression expression) {
+ if (expression.getType() != ExpressionType.FUNCTION) {
+ return;
+ }
+ Function function = expression.getFunctionCall();
+ String functionName = function.getOperator().replace("_",
"").toLowerCase(Locale.ROOT);
+ checkArgument(!CROSS_TABLE_FUNCTIONS.contains(functionName),
+ "Unsupported WHERE clause in DELETE, %s reads another table",
function.getOperator());
+ if (function.getOperands() != null) {
+ for (Expression operand : function.getOperands()) {
+ checkNoCrossTableFunction(operand);
+ }
+ }
+ }
+
+ /// Returns the statement with its table name resolved: qualified with the
database of the request, and in the case
+ /// the table is defined with (table names are case-insensitive by default).
The broker and the controller authorize
+ /// the caller for the resolved table name, and hand the resolved statement
to the executor.
+ ///
+ /// The database of the request is the `database` request header, else the
`database` option of the statement, as
+ /// for a multi-stage query (see
`DatabaseUtils#extractDatabaseFromQueryRequest`). They must match when both are
set,
+ /// and the database of a `database.table` name must match them.
+ ///
+ /// @param databaseHeader value of the `database` request header, if any
+ /// @param tableCache tables of the cluster, to resolve the case of the
table name. A table name it does not know
+ /// keeps the case of the statement.
+ /// @throws QueryException with [QueryErrorCode#QUERY_VALIDATION] if the
`database` header, the `database` option
+ /// and the database of a `database.table` name do
not match (a
+ /// `DatabaseConflictException`), or if the table is
a logical table, which `DELETE` does
+ /// not support
+ public DeleteStatement resolveTableName(@Nullable String databaseHeader,
TableCache tableCache)
+ throws QueryException {
+ String database =
DatabaseUtils.extractDatabaseFromOptionAndHeader(_database, databaseHeader);
+ String tableName;
+ try {
+ tableName = DatabaseUtils.translateTableName(_tableName, database,
tableCache.isIgnoreCase());
+ } catch (IllegalArgumentException e) {
+ throw QueryErrorCode.QUERY_VALIDATION.asException("Invalid table name in
DELETE: " + e.getMessage(), e);
+ }
+ String actualTableName = tableCache.getActualTableName(tableName);
+ if (actualTableName == null &&
tableCache.getActualLogicalTableName(tableName) != null) {
+ // Deleting from a logical table would delete from physical tables the
caller is not authorized for
+ throw QueryErrorCode.QUERY_VALIDATION.asException("DELETE does not
support logical tables: " + tableName);
+ }
+ return new DeleteStatement(actualTableName != null ? actualTableName :
tableName, _predicate, null, _options, true);
+ }
+
+ /// Whether the table name is resolved with [#resolveTableName]. The
executor only executes a resolved statement.
+ public boolean isResolved() {
+ return _resolved;
+ }
+
+ /// Table name: as written in the statement (`table` or `database.table`,
optionally with a type suffix), or, once
+ /// resolved with [#resolveTableName], qualified with its database and in
the case the table is defined with.
+ ///
+ /// The statement handed to the executor is resolved, and the caller is
authorized to delete rows from this table.
+ /// Executors delete from this exact table: resolving the name again, e.g.
case-insensitively, could delete from
+ /// another table than the one the caller is authorized for.
+ public String getTableName() {
+ return _tableName;
+ }
+
+ /// WHERE clause without the `WHERE` keyword, in Pinot SQL that
[CalciteSqlParser#compileToExpression(String)]
+ /// parses into the same expression as the statement. Identifiers are quoted
where the statement quoted them, and
+ /// where needed to keep their name.
+ ///
+ /// It is a standalone expression, e.g. `a = 1 OR b = 2`: wrap it in
parentheses to combine it with other conditions.
+ /// Pinot parses it as an expression but does not validate it as a filter of
the table, e.g. that its columns exist
+ /// or that it has no aggregation, so executors must validate it before
deleting rows.
+ public String getPredicate() {
+ return _predicate;
+ }
+
+ /// Options of the statement other than the database, with the keys as
written: the `SET` statements, the legacy
+ /// `OPTION(...)` suffix and the request `queryOptions`, `SET` taking
precedence over request options with the same
+ /// key. Query options such as `timeoutMs` or `useMultistageEngine` are
included: executors read the options they
+ /// support and ignore the others.
+ ///
+ /// Keys are case-sensitive here, so the same option may appear with
different cases, e.g. `dryRun` set with `SET`
+ /// and `dryrun` in the request: executors that read options
case-insensitively should reject such duplicates rather
+ /// than pick one.
+ public Map<String, String> getOptions() {
+ return _options;
+ }
+
+ @Override
+ public ExecutionType getExecutionType() {
+ return ExecutionType.HTTP;
Review Comment:
Returning ExecutionType.HTTP here doesn't match what HTTP means to the
executor: the HTTP branch of executeStatement calls getResultSchema() and
execute(), and both throw here, so the only thing keeping a DELETE on the right
path is the instanceof DeleteStatement guard that runs before the switch. That
matters a little more now that createSqlQueryExecutor() is a protected hook on
both starters: a subclass that overrides executeStatement and dispatches on the
execution type the way the pre-PR executeDMLStatement body did lands in case
HTTP, calls getResultSchema(), and gets a QUERY_EXECUTION error claiming this
cluster can't delete rows, with its own executeDelete never called. I'd rather
have the enum describe this kind of statement than lean on the concrete type:
add a value like EXECUTOR that DeleteStatement returns, route it to
executeDelete in the switch (keeping the isResolved() check inside that branch)
and drop the instanceof, so a later executor-run DML such as UPDATE pick
s the same value instead of needing another instanceof branch. If you'd prefer
to keep the enum as is, then at least say on getExecutionType() that a
DeleteStatement only runs through SqlQueryExecutor#executeDelete, and have
these three methods throw a message saying they don't apply to DELETE rather
than reusing NOT_SUPPORTED_MESSAGE, which reads as if the cluster rejected the
statement.
##########
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");
+
+ private final String _tableName;
+ private final String _predicate;
+ @Nullable
+ private final String _database;
+ private final Map<String, String> _options;
+ private final boolean _resolved;
+
+ private DeleteStatement(String tableName, String predicate, @Nullable String
database, Map<String, String> options) {
+ this(tableName, predicate, database, options, false);
+ }
+
+ private DeleteStatement(String tableName, String predicate, @Nullable String
database, Map<String, String> options,
+ boolean resolved) {
+ _tableName = tableName;
+ _predicate = predicate;
+ _database = database;
+ _options = Collections.unmodifiableMap(new HashMap<>(options));
+ _resolved = resolved;
+ }
+
+ /// Parses a `DELETE` statement.
+ ///
+ /// @throws IllegalArgumentException if the statement is not a supported
`DELETE`: it must have a WHERE clause, no
+ /// table alias, a WHERE clause that Pinot
can parse as an expression and that does
+ /// not read another table (e.g. with
`lookUp` or `IN_SUBQUERY`), and set the
+ /// database with the `database` option only
+ public static DeleteStatement parse(SqlNodeAndOptions sqlNodeAndOptions) {
+ SqlNode sqlNode = sqlNodeAndOptions.getSqlNode();
+ checkArgument(sqlNode instanceof SqlDelete, "Not a DELETE statement: %s",
sqlNode.getKind());
+ SqlDelete sqlDelete = (SqlDelete) sqlNode;
+ String tableName = getTableName(sqlDelete.getTargetTable());
+ checkArgument(sqlDelete.getAlias() == null,
+ "DELETE does not support a table alias, reference the columns
directly");
+ SqlNode condition = sqlDelete.getCondition();
+ checkArgument(condition != null,
+ "DELETE requires a WHERE clause; delete the table segments to remove
all of its rows");
+ String predicate = toPinotSql(condition);
+
+ String database = null;
+ Map<String, String> options = new HashMap<>();
+ for (Map.Entry<String, String> option :
sqlNodeAndOptions.getOptions().entrySet()) {
+ String key = option.getKey();
+ if (key.equals(CommonConstants.DATABASE)) {
+ database = option.getValue();
+ } else if (key.equalsIgnoreCase(CommonConstants.DATABASE)) {
+ // Queries ignore it, as they only read the `database` option: fail
rather than delete from another table
+ throw new IllegalArgumentException(
+ "Unsupported option: " + key + ", set the database with the '" +
CommonConstants.DATABASE + "' option");
+ } else {
+ options.put(key, option.getValue());
+ }
+ }
+ return new DeleteStatement(tableName, predicate, database, options);
+ }
+
+ private static String getTableName(SqlNode targetTable) {
+ // Table hints and EXTEND clauses parse into other node types
+ checkArgument(targetTable instanceof SqlIdentifier, "DELETE only supports
a plain table name, got: %s",
+ targetTable);
+ // A quoted name part may contain a dot, which splits it like the table
name of a query
+ String tableName = String.join(".", ((SqlIdentifier) targetTable).names);
+ // Empty parts, e.g. in "db."."t", are rejected rather than dropped, so
that the table name is the one authorized
+ String[] parts = tableName.split("\\.", -1);
+ checkArgument(parts.length <= 2 &&
Arrays.stream(parts).noneMatch(String::isEmpty),
+ "Invalid table name: %s, expected [database.]table", tableName);
+ return tableName;
+ }
+
+ /// Serializes a WHERE clause back into SQL, and verifies that Pinot parses
it into the same expression.
+ private static String toPinotSql(SqlNode condition) {
+ Expression expression;
+ Expression serializedExpression;
+ String predicate;
+ try {
+ expression = CalciteSqlParser.compileToExpression(condition);
+ predicate = condition.accept(KEYWORD_IDENTIFIER_QUOTER)
+ .toSqlString(config -> config.withDialect(PINOT_SQL_DIALECT)
+ .withQuoteAllIdentifiers(false)
+ .withIndentation(0))
+ .getSql();
+ serializedExpression = CalciteSqlParser.compileToExpression(predicate);
+ } catch (Exception e) {
+ throw new IllegalArgumentException("Unsupported WHERE clause in DELETE:
" + condition, e);
+ }
+ // Fail rather than hand over a predicate that selects other rows than the
statement
+ checkArgument(serializedExpression.equals(expression),
+ "Unsupported WHERE clause in DELETE, it cannot be serialized back into
the same expression: %s", condition);
+ checkNoCrossTableFunction(expression);
+ return predicate;
+ }
+
+ /// Rejects the functions that read another table: the caller is only
authorized for the table it deletes from.
+ private static void checkNoCrossTableFunction(Expression expression) {
+ if (expression.getType() != ExpressionType.FUNCTION) {
+ return;
+ }
+ Function function = expression.getFunctionCall();
+ String functionName = function.getOperator().replace("_",
"").toLowerCase(Locale.ROOT);
+ checkArgument(!CROSS_TABLE_FUNCTIONS.contains(functionName),
+ "Unsupported WHERE clause in DELETE, %s reads another table",
function.getOperator());
+ if (function.getOperands() != null) {
+ for (Expression operand : function.getOperands()) {
+ checkNoCrossTableFunction(operand);
+ }
+ }
+ }
+
+ /// Returns the statement with its table name resolved: qualified with the
database of the request, and in the case
+ /// the table is defined with (table names are case-insensitive by default).
The broker and the controller authorize
+ /// the caller for the resolved table name, and hand the resolved statement
to the executor.
+ ///
+ /// The database of the request is the `database` request header, else the
`database` option of the statement, as
+ /// for a multi-stage query (see
`DatabaseUtils#extractDatabaseFromQueryRequest`). They must match when both are
set,
+ /// and the database of a `database.table` name must match them.
+ ///
+ /// @param databaseHeader value of the `database` request header, if any
+ /// @param tableCache tables of the cluster, to resolve the case of the
table name. A table name it does not know
+ /// keeps the case of the statement.
+ /// @throws QueryException with [QueryErrorCode#QUERY_VALIDATION] if the
`database` header, the `database` option
+ /// and the database of a `database.table` name do
not match (a
+ /// `DatabaseConflictException`), or if the table is
a logical table, which `DELETE` does
+ /// not support
+ public DeleteStatement resolveTableName(@Nullable String databaseHeader,
TableCache tableCache)
+ throws QueryException {
+ String database =
DatabaseUtils.extractDatabaseFromOptionAndHeader(_database, databaseHeader);
+ String tableName;
+ try {
+ tableName = DatabaseUtils.translateTableName(_tableName, database,
tableCache.isIgnoreCase());
+ } catch (IllegalArgumentException e) {
+ throw QueryErrorCode.QUERY_VALIDATION.asException("Invalid table name in
DELETE: " + e.getMessage(), e);
+ }
+ String actualTableName = tableCache.getActualTableName(tableName);
+ if (actualTableName == null &&
tableCache.getActualLogicalTableName(tableName) != null) {
+ // Deleting from a logical table would delete from physical tables the
caller is not authorized for
+ throw QueryErrorCode.QUERY_VALIDATION.asException("DELETE does not
support logical tables: " + tableName);
+ }
+ return new DeleteStatement(actualTableName != null ? actualTableName :
tableName, _predicate, null, _options, true);
+ }
+
+ /// Whether the table name is resolved with [#resolveTableName]. The
executor only executes a resolved statement.
+ public boolean isResolved() {
+ return _resolved;
+ }
+
+ /// Table name: as written in the statement (`table` or `database.table`,
optionally with a type suffix), or, once
+ /// resolved with [#resolveTableName], qualified with its database and in
the case the table is defined with.
+ ///
+ /// The statement handed to the executor is resolved, and the caller is
authorized to delete rows from this table.
+ /// Executors delete from this exact table: resolving the name again, e.g.
case-insensitively, could delete from
+ /// another table than the one the caller is authorized for.
+ public String getTableName() {
+ return _tableName;
+ }
+
+ /// WHERE clause without the `WHERE` keyword, in Pinot SQL that
[CalciteSqlParser#compileToExpression(String)]
+ /// parses into the same expression as the statement. Identifiers are quoted
where the statement quoted them, and
+ /// where needed to keep their name.
+ ///
+ /// It is a standalone expression, e.g. `a = 1 OR b = 2`: wrap it in
parentheses to combine it with other conditions.
+ /// Pinot parses it as an expression but does not validate it as a filter of
the table, e.g. that its columns exist
+ /// or that it has no aggregation, so executors must validate it before
deleting rows.
+ public String getPredicate() {
+ return _predicate;
+ }
+
+ /// Options of the statement other than the database, with the keys as
written: the `SET` statements, the legacy
Review Comment:
Small doc mismatch, but it's on the new executor contract so worth getting
right: this Javadoc (and the class Javadoc above, "reaches the executor as
written") says keys arrive as written, yet every option source is run through
`QueryOptionsUtils.resolveCaseInsensitiveOptions` before `parse` sees it (SET
options in `CalciteSqlParser.compileToSqlNodeAndOptions`, legacy `OPTION(...)`
in `SqlNodeAndOptions.setExtraOptions`, request `queryOptions` in
`RequestUtils.setOptions`), and `parse` copies the map through unchanged. So
`SET TIMEOUTMS = 1000; DELETE FROM t WHERE a = 1` surfaces as
`{timeoutMs=1000}`, and `SET timeoutms = 1` plus request
`queryOptions=TIMEOUTMS=2` collapses into a single `timeoutMs=1` entry rather
than the two-key duplicate the second paragraph tells executors to detect. Your
own test already says this ("with their canonical keys" in
`DeleteStatementTest`), so it's just the prose that is off; an implementer who
trusts it and reads `getOptions().get("timeoutms")`
gets null and silently falls back to a default. I'd reword it to say keys are
kept as written except that names declared on
`CommonConstants.Broker.Request.QueryOptionKey` are canonicalized
case-insensitively as for queries (so those never appear twice), while
executor-specific keys like `dryRun` and keys registered with
`registerSqlQueryOptionKey` keep their case and can still show up under two
spellings. The class-level sentence could just become "is passed to the
executor through [#getOptions()]", since its point is that nothing gets
dropped, not key spelling.
##########
pinot-broker/src/main/java/org/apache/pinot/broker/api/resources/PinotClientRequest.java:
##########
@@ -713,6 +738,112 @@ private BrokerResponse executeSqlQuery(ObjectNode
sqlRequestJson, HttpRequesterI
}
}
+ /// Executes a `DELETE` once the caller is authorized to delete rows from
its table.
+ ///
+ /// The first-step access control runs first, as for queries. The table is
then resolved with the database of the
+ /// request and in the case it is defined with, the caller is authorized to
delete rows from it (see
+ /// [#authorizeDelete]), and the executor deletes rows from that exact table.
+ private BrokerResponse executeDelete(SqlNodeAndOptions sqlNodeAndOptions,
Map<String, String> headers,
+ HttpRequesterIdentity requesterIdentity, @Nullable HttpHeaders
httpHeaders) {
+ AccessControl accessControl = _accessControlFactory.create();
+ // The first-step access control runs before the table is looked up, as
for queries
+ AuthorizationResult authorizationResult =
accessControl.authorize(requesterIdentity);
+ if (!authorizationResult.hasAccess()) {
+ throw deleteAccessDenied(null, authorizationResult);
+ }
+ if (_tableCache == null) {
+ return new BrokerResponseNative(QueryErrorCode.QUERY_VALIDATION,
+ "DELETE is not supported by this broker: no table cache was
configured");
+ }
+ DeleteStatement statement;
+ try {
+ String databaseHeader = httpHeaders != null ?
httpHeaders.getHeaderString(CommonConstants.DATABASE) : null;
+ statement = ((DeleteStatement)
DataManipulationStatementParser.parse(sqlNodeAndOptions))
+ .resolveTableName(databaseHeader, _tableCache);
+ } catch (QueryException e) {
+ // e.g. an invalid statement, a logical table, or a database header that
does not match the statement
+ return new BrokerResponseNative(e.getErrorCode(), e.getMessage());
+ }
+ authorizeDelete(accessControl, statement.getTableName(),
requesterIdentity, httpHeaders);
+ return _sqlQueryExecutor.executeStatement(statement, headers);
Review Comment:
One thing that stands out on this path is the asymmetry: a denied DELETE
gets an INFO line plus REQUEST_DROPPED_DUE_TO_ACCESS_ERROR in
deleteAccessDenied, but a DELETE that passes authorization goes straight into
the executor here with no log line or meter at all. SqlQueryExecutor has no
logger and its Javadoc explicitly says it won't log the statement as a query,
BrokerRequestHandlerDelegate rejects SqlDelete so it never reaches the broker
query log, and the Jersey audit filter is off by default, so in a deployment
that overrides executeDelete the only record of a row deletion is whatever that
override chooses to write. For something that removes data I'd rather the entry
point always leave a trace, the way segment and schema deletions log at INFO,
and the same applies to PinotQueryResource.executeDelete on the controller.
Could we log the resolved table, predicate, option keys and requester right
before the executor call (either in both entry points, or once in
executeStatement
so every override inherits it), and bump a counter so attempts show up next to
the denial meter? Something like:
```java
LOGGER.info("Executing DELETE on table: {}, predicate: {}, options: {},
client: {}", statement.getTableName(),
statement.getPredicate(), statement.getOptions().keySet(),
requesterIdentity.getClientIp());
```
##########
pinot-controller/src/main/java/org/apache/pinot/controller/BaseControllerStarter.java:
##########
@@ -669,7 +669,7 @@ private void setUpPinotController() {
new SegmentCompletionManager(_helixParticipantManager,
_pinotLLCRealtimeSegmentManager, _controllerMetrics,
_leadControllerManager, _config.getSegmentCommitTimeoutSeconds(),
segmentCompletionConfig);
- _sqlQueryExecutor = new SqlQueryExecutor(_config.generateVipUrl());
+ _sqlQueryExecutor = createSqlQueryExecutor();
Review Comment:
This call runs before `_connectionManager` is built two lines down and well
before `setupControllerPeriodicTasks()` creates `_taskManager`, and now that
it's an overridable hook that ordering quietly becomes part of the contract.
The broker-side `createSqlQueryExecutor()` is invoked after the request handler
and table cache are up, so an override written the same way here, say `new
PurgeSqlQueryExecutor(_config.generateVipUrl(), _taskManager)` to create ad-hoc
purge tasks the way `PinotTaskRestletResource.executeAdhocTask` does, captures
null and the first authorized `DELETE` on `/sql` NPEs inside the
`StreamingOutput` and surfaces as a 500. Nothing reads `_sqlQueryExecutor`
until the HK2 binder, so it should be safe to move this line down to just
before `_adminApp.registerBinder(...)`, after the periodic tasks are set up,
which gives the hook the same fully-initialized view the broker hook gets. If
you'd rather keep it here, it would help to state in the hook's Javadoc which
comp
onents are available at this point (`_config`, `_helixResourceManager`,
`_helixTaskResourceManager`, `_pinotLLCRealtimeSegmentManager`) and that
`_taskManager` and `_connectionManager` are not yet created.
--
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]