This is an automated email from the ASF dual-hosted git repository.

chia7712 pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git


The following commit(s) were added to refs/heads/trunk by this push:
     new 17f3f8aeae8 MINOR: Cleanups following KAFKA-20637 (#22584)
17f3f8aeae8 is described below

commit 17f3f8aeae8c7a93fce731b016ffb055c2bee4bc
Author: Mickael Maison <[email protected]>
AuthorDate: Wed Jul 29 04:56:57 2026 +0200

    MINOR: Cleanups following KAFKA-20637 (#22584)
    
    Following up comment on https://github.com/apache/kafka/pull/22408
    
    Reviewers: Chia-Ping Tsai <[email protected]>
---
 .../main/scala/kafka/server/ControllerApis.scala   |  2 +-
 .../scala/kafka/server/RequestHandlerHelper.scala  |  4 +--
 .../scala/kafka/tools/TestRaftRequestHandler.scala |  6 ++--
 .../unit/kafka/network/SocketServerTest.scala      |  4 +--
 .../java/org/apache/kafka/network/Request.java     |  4 +++
 .../org/apache/kafka/server/ForwardingManager.java | 32 ++++------------------
 .../apache/kafka/server/ForwardingManagerImpl.java | 25 ++++++++---------
 .../org/apache/kafka/server/EnvelopeUtilsTest.java |  2 +-
 8 files changed, 29 insertions(+), 50 deletions(-)

diff --git a/core/src/main/scala/kafka/server/ControllerApis.scala 
b/core/src/main/scala/kafka/server/ControllerApis.scala
index 9f56906f822..cc42cba56b8 100644
--- a/core/src/main/scala/kafka/server/ControllerApis.scala
+++ b/core/src/main/scala/kafka/server/ControllerApis.scala
@@ -686,7 +686,7 @@ class ControllerApis(
     request: Request,
     buildResponse: ApiMessage => AbstractResponse
   ): CompletableFuture[Unit] = {
-    val requestBody = request.body(classOf[AbstractRequest])
+    val requestBody = request.body
     val future = raftManager
       .handleRequest(request.context, request.header, requestBody.data, 
time.milliseconds())
       .toCompletableFuture
diff --git a/core/src/main/scala/kafka/server/RequestHandlerHelper.scala 
b/core/src/main/scala/kafka/server/RequestHandlerHelper.scala
index 29ee72e025b..9c78451d674 100644
--- a/core/src/main/scala/kafka/server/RequestHandlerHelper.scala
+++ b/core/src/main/scala/kafka/server/RequestHandlerHelper.scala
@@ -20,7 +20,7 @@ package kafka.server
 import kafka.network.RequestChannel
 import kafka.server.QuotaFactory.QuotaManagers
 import org.apache.kafka.common.errors.ClusterAuthorizationException
-import org.apache.kafka.common.requests.{AbstractRequest, AbstractResponse}
+import org.apache.kafka.common.requests.AbstractResponse
 import org.apache.kafka.common.utils.Time
 import org.apache.kafka.network.Request
 import org.apache.kafka.server.quota.{ClientQuotaManager, 
ControllerMutationQuota, ThrottleCallback}
@@ -56,7 +56,7 @@ class RequestHandlerHelper(
     error: Throwable,
     throttleMs: Int
   ): Unit = {
-    val requestBody = request.body(classOf[AbstractRequest])
+    val requestBody = request.body
     val response = requestBody.getErrorResponse(throttleMs, error)
     if (response == null)
       requestChannel.closeConnection(request, requestBody.errorCounts(error))
diff --git a/core/src/main/scala/kafka/tools/TestRaftRequestHandler.scala 
b/core/src/main/scala/kafka/tools/TestRaftRequestHandler.scala
index 28badb26c3f..d9646e8dff9 100644
--- a/core/src/main/scala/kafka/tools/TestRaftRequestHandler.scala
+++ b/core/src/main/scala/kafka/tools/TestRaftRequestHandler.scala
@@ -23,7 +23,7 @@ import kafka.utils.Logging
 import org.apache.kafka.common.internals.FatalExitError
 import org.apache.kafka.common.message.{BeginQuorumEpochResponseData, 
EndQuorumEpochResponseData, FetchResponseData, FetchSnapshotResponseData, 
VoteResponseData}
 import org.apache.kafka.common.protocol.{ApiKeys, ApiMessage}
-import org.apache.kafka.common.requests.{AbstractRequest, AbstractResponse, 
BeginQuorumEpochResponse, EndQuorumEpochResponse, FetchResponse, 
FetchSnapshotResponse, VoteResponse}
+import org.apache.kafka.common.requests.{AbstractResponse, 
BeginQuorumEpochResponse, EndQuorumEpochResponse, FetchResponse, 
FetchSnapshotResponse, VoteResponse}
 import org.apache.kafka.common.utils.Time
 import org.apache.kafka.network.Request
 import org.apache.kafka.raft.RaftManager
@@ -57,7 +57,7 @@ class TestRaftRequestHandler(
       case e: Throwable =>
         error(s"Unexpected error handling request ${request.requestDesc(true)} 
" +
           s"with context ${request.context}", e)
-        val errorResponse = 
request.body(classOf[AbstractRequest]).getErrorResponse(e)
+        val errorResponse = request.body.getErrorResponse(e)
         requestChannel.sendResponse(request, errorResponse)
     } finally {
       // The local completion time may be set while processing the request. 
Only record it if it's unset.
@@ -94,7 +94,7 @@ class TestRaftRequestHandler(
     request: Request,
     buildResponse: ApiMessage => AbstractResponse
   ): Unit = {
-    val requestBody = request.body(classOf[AbstractRequest])
+    val requestBody = request.body
 
     val future = raftManager.handleRequest(
       request.context,
diff --git a/core/src/test/scala/unit/kafka/network/SocketServerTest.scala 
b/core/src/test/scala/unit/kafka/network/SocketServerTest.scala
index 51ea50628ef..bc4cfe4a513 100644
--- a/core/src/test/scala/unit/kafka/network/SocketServerTest.scala
+++ b/core/src/test/scala/unit/kafka/network/SocketServerTest.scala
@@ -161,7 +161,7 @@ class SocketServerTest {
   }
 
   def processRequest(channel: RequestChannel, request: Request): Unit = {
-    val byteBuffer = 
request.body(classOf[AbstractRequest]).serializeWithHeader(request.header)
+    val byteBuffer = request.body.serializeWithHeader(request.header)
     val send = new NetworkSend(request.context.connectionId, 
ByteBufferSend.sizePrefixed(byteBuffer))
     val headerLog = RequestConvertToJson.requestHeaderNode(request.header)
     channel.sendResponse(new SendResponse(request, send, 
Optional.of(headerLog)))
@@ -641,7 +641,7 @@ class SocketServerTest {
     // Mimic a primitive request handler that fetches the request from 
RequestChannel and place a response with a
     // throttled channel.
     val request = receiveRequest(server.dataPlaneRequestChannel)
-    val byteBuffer = 
request.body(classOf[AbstractRequest]).serializeWithHeader(request.header)
+    val byteBuffer = request.body.serializeWithHeader(request.header)
     val send = new NetworkSend(request.context.connectionId, 
ByteBufferSend.sizePrefixed(byteBuffer))
 
     val channelThrottlingCallback = new ThrottleCallback {
diff --git a/server/src/main/java/org/apache/kafka/network/Request.java 
b/server/src/main/java/org/apache/kafka/network/Request.java
index c9144889496..e9258e4e3a1 100644
--- a/server/src/main/java/org/apache/kafka/network/Request.java
+++ b/server/src/main/java/org/apache/kafka/network/Request.java
@@ -282,6 +282,10 @@ public final class Request implements BaseRequest {
             + ", but found " + bodyAndSize.request.getClass().getName());
     }
 
+    public AbstractRequest body() {
+        return bodyAndSize.request;
+    }
+
     public AbstractRequest loggableRequest() {
         if (bodyAndSize.request instanceof AlterConfigsRequest alterConfigs) {
             var newData = alterConfigs.data().duplicate();
diff --git 
a/server/src/main/java/org/apache/kafka/server/ForwardingManager.java 
b/server/src/main/java/org/apache/kafka/server/ForwardingManager.java
index 79e26566929..708e137c3a6 100644
--- a/server/src/main/java/org/apache/kafka/server/ForwardingManager.java
+++ b/server/src/main/java/org/apache/kafka/server/ForwardingManager.java
@@ -19,20 +19,13 @@ package org.apache.kafka.server;
 import org.apache.kafka.clients.NodeApiVersions;
 import org.apache.kafka.common.requests.AbstractRequest;
 import org.apache.kafka.common.requests.AbstractResponse;
-import org.apache.kafka.common.requests.RequestContext;
 import org.apache.kafka.network.Request;
 
 import java.nio.ByteBuffer;
 import java.util.Optional;
 import java.util.function.Consumer;
-import java.util.function.Supplier;
 
-public interface ForwardingManager {
-
-    /**
-     * Close the forwarding manager
-     */
-    void close();
+public interface ForwardingManager extends AutoCloseable {
 
     /**
      * Forward given request to the active controller.
@@ -49,12 +42,7 @@ public interface ForwardingManager {
     ) {
         ByteBuffer buffer = originalRequest.buffer().duplicate();
         buffer.flip();
-        forwardRequest(originalRequest.context(),
-                buffer,
-                originalRequest.startTimeNanos(),
-                originalRequest.body(AbstractRequest.class),
-                originalRequest::toString,
-                responseCallback);
+        forwardRequest(originalRequest, buffer, originalRequest.body(), 
responseCallback);
     }
 
     /**
@@ -72,36 +60,26 @@ public interface ForwardingManager {
             AbstractRequest newRequestBody,
             Consumer<Optional<AbstractResponse>> responseCallback) {
         ByteBuffer buffer = 
newRequestBody.serializeWithHeader(originalRequest.header());
-        forwardRequest(originalRequest.context(),
-                buffer,
-                originalRequest.startTimeNanos(),
-                newRequestBody,
-                originalRequest::toString,
-                responseCallback);
+        forwardRequest(originalRequest, buffer, newRequestBody, 
responseCallback);
     }
 
     /**
      * Forward given request to the active controller.
      *
-     * @param requestContext      The request context of the original envelope 
request.
+     * @param originalRequest     The request to forward.
      * @param requestBufferCopy   The request buffer we want to send. This 
should not be the original
      *                            byte buffer from the envelope request, since 
we will be mutating
      *                            the position and limit fields. It should be 
a copy.
-     * @param requestCreationNs   The request creation timestamp in 
nanoseconds.
      * @param requestBody         The AbstractRequest we are sending.
-     * @param requestToString     A callback which can be invoked to produce a 
human-readable description
-     *                            of the request.
      * @param responseCallback    A callback which takes in an 
`Optional&lt;AbstractResponse&gt;`.
      *                            We will call this function with 
Optional.of(x) after the controller responds with x.
      *                            Or, if the controller doesn't support the 
request version, we will complete
      *                            the callback with Optional.empty().
      */
     void forwardRequest(
-            RequestContext requestContext,
+            Request originalRequest,
             ByteBuffer requestBufferCopy,
-            long requestCreationNs,
             AbstractRequest requestBody,
-            Supplier<String> requestToString,
             Consumer<Optional<AbstractResponse>> responseCallback);
 
     /**
diff --git 
a/server/src/main/java/org/apache/kafka/server/ForwardingManagerImpl.java 
b/server/src/main/java/org/apache/kafka/server/ForwardingManagerImpl.java
index 0344bd67d74..63cc4b54572 100644
--- a/server/src/main/java/org/apache/kafka/server/ForwardingManagerImpl.java
+++ b/server/src/main/java/org/apache/kafka/server/ForwardingManagerImpl.java
@@ -25,8 +25,8 @@ import org.apache.kafka.common.requests.AbstractRequest;
 import org.apache.kafka.common.requests.AbstractResponse;
 import org.apache.kafka.common.requests.EnvelopeRequest;
 import org.apache.kafka.common.requests.EnvelopeResponse;
-import org.apache.kafka.common.requests.RequestContext;
 import org.apache.kafka.common.requests.RequestHeader;
+import org.apache.kafka.network.Request;
 import org.apache.kafka.server.common.ControllerRequestCompletionHandler;
 import org.apache.kafka.server.common.NodeToControllerChannelManager;
 import org.apache.kafka.server.metrics.ForwardingManagerMetrics;
@@ -38,9 +38,8 @@ import java.nio.ByteBuffer;
 import java.util.Optional;
 import java.util.concurrent.TimeUnit;
 import java.util.function.Consumer;
-import java.util.function.Supplier;
 
-public class ForwardingManagerImpl implements ForwardingManager, AutoCloseable 
{
+public class ForwardingManagerImpl implements ForwardingManager {
 
     private static final Logger LOG = 
LoggerFactory.getLogger(ForwardingManagerImpl.class);
 
@@ -58,14 +57,12 @@ public class ForwardingManagerImpl implements 
ForwardingManager, AutoCloseable {
 
     @Override
     public void forwardRequest(
-            RequestContext requestContext,
+            Request originalRequest,
             ByteBuffer requestBufferCopy,
-            long requestCreationNs,
             AbstractRequest requestBody,
-            Supplier<String> requestToString,
             Consumer<Optional<AbstractResponse>> responseCallback) {
-        EnvelopeRequest.Builder envelopeRequest = 
ForwardingManagerUtil.buildEnvelopeRequest(requestContext, requestBufferCopy);
-        long requestCreationTimeMs = 
TimeUnit.NANOSECONDS.toMillis(requestCreationNs);
+        EnvelopeRequest.Builder envelopeRequest = 
ForwardingManagerUtil.buildEnvelopeRequest(originalRequest.context(), 
requestBufferCopy);
+        long requestCreationTimeMs = 
TimeUnit.NANOSECONDS.toMillis(originalRequest.startTimeNanos());
 
         class ForwardingResponseHandler implements 
ControllerRequestCompletionHandler {
 
@@ -77,11 +74,11 @@ public class ForwardingManagerImpl implements 
ForwardingManager, AutoCloseable {
 
                 if (clientResponse.versionMismatch() != null) {
                     LOG.debug("Returning `UNKNOWN_SERVER_ERROR` in response to 
{} due to unexpected version error",
-                            requestToString.get(), 
clientResponse.versionMismatch());
+                            originalRequest, clientResponse.versionMismatch());
                     
responseCallback.accept(Optional.of(requestBody.getErrorResponse(Errors.UNKNOWN_SERVER_ERROR.exception())));
                 } else if (clientResponse.authenticationException() != null) {
                     LOG.debug("Returning `UNKNOWN_SERVER_ERROR` in response to 
{} due to authentication error",
-                            requestToString.get(), 
clientResponse.authenticationException());
+                            originalRequest, 
clientResponse.authenticationException());
                     
responseCallback.accept(Optional.of(requestBody.getErrorResponse(Errors.UNKNOWN_SERVER_ERROR.exception())));
                 } else {
                     EnvelopeResponse envelopeResponse = (EnvelopeResponse) 
clientResponse.responseBody();
@@ -102,10 +99,10 @@ public class ForwardingManagerImpl implements 
ForwardingManager, AutoCloseable {
                             // return `UNKNOWN_SERVER_ERROR` so that the user 
knows that there is a problem
                             // on the broker.
                             LOG.debug("Forwarded request {} failed with an 
error in the envelope response {}",
-                                    requestToString.get(), envelopeError);
+                                    originalRequest, envelopeError);
                             response = 
requestBody.getErrorResponse(Errors.UNKNOWN_SERVER_ERROR.exception());
                         } else {
-                            response = 
parseResponse(envelopeResponse.responseData(), requestBody, 
requestContext.header);
+                            response = 
parseResponse(envelopeResponse.responseData(), requestBody, 
originalRequest.context().header);
                         }
                         responseCallback.accept(Optional.of(response));
                     }
@@ -114,7 +111,7 @@ public class ForwardingManagerImpl implements 
ForwardingManager, AutoCloseable {
 
             @Override
             public void onTimeout() {
-                LOG.debug("Forwarding of the request {} failed due to timeout 
exception", requestToString.get());
+                LOG.debug("Forwarding of the request {} failed due to timeout 
exception", originalRequest);
                 forwardingManagerMetrics.decrementQueueLength();
                 
forwardingManagerMetrics.queueTimeMsHist().record(channelManager.getTimeoutMs());
                 AbstractResponse response = requestBody.getErrorResponse(new 
TimeoutException());
@@ -136,7 +133,7 @@ public class ForwardingManagerImpl implements 
ForwardingManager, AutoCloseable {
         return channelManager.controllerApiVersions();
     }
 
-    private AbstractResponse parseResponse(ByteBuffer buffer, AbstractRequest 
request, RequestHeader header) {
+    private static AbstractResponse parseResponse(ByteBuffer buffer, 
AbstractRequest request, RequestHeader header) {
         try {
             return AbstractResponse.parseResponse(buffer, header);
         } catch (Exception e) {
diff --git 
a/server/src/test/java/org/apache/kafka/server/EnvelopeUtilsTest.java 
b/server/src/test/java/org/apache/kafka/server/EnvelopeUtilsTest.java
index e365fab9575..af7e6351bc2 100644
--- a/server/src/test/java/org/apache/kafka/server/EnvelopeUtilsTest.java
+++ b/server/src/test/java/org/apache/kafka/server/EnvelopeUtilsTest.java
@@ -86,7 +86,7 @@ public class EnvelopeUtilsTest {
         assertTrue(forwardedRequest.isForwarded());
         assertSame(envelope, forwardedRequest.envelope().orElseThrow());
         assertEquals(envelope.requestDequeueTimeNanos(), 
forwardedRequest.requestDequeueTimeNanos());
-        assertInstanceOf(CreateTopicsRequest.class, 
forwardedRequest.body(AbstractRequest.class));
+        assertInstanceOf(CreateTopicsRequest.class, forwardedRequest.body());
     }
 
     @Test

Reply via email to