xiangfu0 commented on code in PR #19682:
URL: https://github.com/apache/pinot/pull/19682#discussion_r4119855513
##########
pinot-query-runtime/src/main/java/org/apache/pinot/query/mailbox/GrpcSendingMailbox.java:
##########
@@ -297,7 +297,8 @@ public void cancel(Throwable t) {
String msg = t != null ? t.getMessage() : "Unknown";
// NOTE: DO NOT use onError() because it will terminate the stream, and
receiver might not get the callback
MseBlock errorBlock = ErrorMseBlock.fromError(
- QueryErrorCode.QUERY_CANCELLATION, "Cancelled by sender with
exception: " + msg);
+ QueryErrorCode.fromThrowable(t, QueryErrorCode.QUERY_CANCELLATION),
Review Comment:
Thanks for tracing the race. This PR preserves the error code in the
cancellation fallback; it does not change the scheduler timing.
##########
pinot-query-runtime/src/main/java/org/apache/pinot/query/mailbox/GrpcSendingMailbox.java:
##########
@@ -297,7 +297,8 @@ public void cancel(Throwable t) {
String msg = t != null ? t.getMessage() : "Unknown";
// NOTE: DO NOT use onError() because it will terminate the stream, and
receiver might not get the callback
MseBlock errorBlock = ErrorMseBlock.fromError(
- QueryErrorCode.QUERY_CANCELLATION, "Cancelled by sender with
exception: " + msg);
+ QueryErrorCode.fromThrowable(t, QueryErrorCode.QUERY_CANCELLATION),
+ "Cancelled by sender with exception: " + msg);
Review Comment:
Updated both mailbox paths: non-QUERY_CANCELLATION errors forward the
original message, while genuine cancellations keep the prefix. The focused
tests assert both cases.
##########
pinot-query-runtime/src/main/java/org/apache/pinot/query/mailbox/GrpcSendingMailbox.java:
##########
@@ -134,8 +134,8 @@ public void send(MseBlock.Eos block, List<DataBuffer>
serializedStats) {
// and must always reach the receiver, so they bypass the back-pressure
gate. Bypassing also disables the
// cooperative termination poll inside [#awaitReady]; without it, a
terminate signal raised while the sender is
// mid-way through pushing an error EOS would unwind [#sendInternal] with
a TerminationException, leave
- // [#_senderSideClosed] false, and let [#cancel] run and overwrite the
original error code with
- // QUERY_CANCELLATION on the receiver side.
+ // [#_senderSideClosed] false, and let [#cancel] run and replace the
original error code with
+ // QUERY_CANCELLATION if its throwable does not carry a QueryErrorCode.
Review Comment:
Updated the PR description to state that all QueryException codes, including
EXECUTION_TIMEOUT and INTERNAL, now propagate and change user-visible error
codes and metrics.
##########
pinot-query-runtime/src/main/java/org/apache/pinot/query/mailbox/InMemorySendingMailbox.java:
##########
@@ -120,8 +120,8 @@ public void cancel(Throwable t) {
_receivingMailbox = _mailboxService.getReceivingMailbox(_id);
}
_receivingMailbox.setErrorBlock(
- ErrorMseBlock.fromException(new QueryCancelledException(
- "Cancelled by sender with exception: " + t.getMessage())),
List.of());
+ ErrorMseBlock.fromError(QueryErrorCode.fromThrowable(t,
QueryErrorCode.QUERY_CANCELLATION),
Review Comment:
Updated the in-memory path to handle null throwables and null messages,
matching gRPC, with focused test rows for both transports.
##########
pinot-query-runtime/src/test/java/org/apache/pinot/query/mailbox/GrpcSendingMailboxTest.java:
##########
@@ -50,19 +51,55 @@
import org.apache.pinot.segment.spi.memory.DataBuffer;
import org.apache.pinot.segment.spi.memory.PinotByteBuffer;
import org.apache.pinot.spi.exception.QueryErrorCode;
+import org.apache.pinot.spi.exception.QueryException;
import org.apache.pinot.spi.exception.TerminationException;
import org.apache.pinot.spi.query.QueryThreadContext;
+import org.mockito.ArgumentCaptor;
import org.mockito.Mockito;
import org.testng.Assert;
import org.testng.annotations.DataProvider;
import org.testng.annotations.Test;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
import static org.testng.Assert.assertEquals;
import static org.testng.Assert.assertTrue;
public class GrpcSendingMailboxTest {
+ @Test(dataProvider = "cancellationErrors")
+ @SuppressWarnings("unchecked")
+ public void cancelPreservesQueryErrorCode(Exception exception,
QueryErrorCode expectedCode)
+ throws IOException {
Review Comment:
Added TerminationException(EXECUTION_TIMEOUT) and null throwable rows in
both mailbox tests. The final focused JDK 25 run passed 34 tests.
--
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]