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]

Reply via email to