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 faee85e4e54 MINOR: Move the ForwardingManager helper to the server 
module (#21331)
faee85e4e54 is described below

commit faee85e4e54caa392af4246e09153910c4b6e0eb
Author: Lan Ding <[email protected]>
AuthorDate: Wed Jan 21 00:41:23 2026 +0800

    MINOR: Move the ForwardingManager helper to the server module (#21331)
    
    Move the `ForwardingManager#buildEnvelopeRequest` helper to the server
    module.
    
    Reviewers: Chia-Ping Tsai <[email protected]>
---
 .../kafka/server/AutoTopicCreationManager.scala    |  1 +
 .../src/main/scala/kafka/server/BrokerServer.scala |  6 ++--
 .../scala/kafka/server/ForwardingManager.scala     | 29 ++-------------
 .../server/AutoTopicCreationManagerTest.scala      |  1 +
 .../apache/kafka/server/ForwardingManagerUtil.java | 41 ++++++++++++++++++++++
 .../apache}/kafka/server/KRaftTopicCreator.java    |  4 +--
 .../org/apache}/kafka/server/TopicCreator.java     |  2 +-
 .../kafka/server/KRaftTopicCreatorTest.java        |  2 +-
 8 files changed, 52 insertions(+), 34 deletions(-)

diff --git a/core/src/main/scala/kafka/server/AutoTopicCreationManager.scala 
b/core/src/main/scala/kafka/server/AutoTopicCreationManager.scala
index f66aa7d3c1a..6f2a192a438 100644
--- a/core/src/main/scala/kafka/server/AutoTopicCreationManager.scala
+++ b/core/src/main/scala/kafka/server/AutoTopicCreationManager.scala
@@ -35,6 +35,7 @@ import org.apache.kafka.coordinator.share.ShareCoordinator
 import org.apache.kafka.coordinator.transaction.TransactionLogConfig
 import org.apache.kafka.server.quota.ControllerMutationQuota
 import org.apache.kafka.common.utils.Time
+import org.apache.kafka.server.TopicCreator
 
 import scala.collection.{Map, Seq, Set, mutable}
 import scala.jdk.CollectionConverters._
diff --git a/core/src/main/scala/kafka/server/BrokerServer.scala 
b/core/src/main/scala/kafka/server/BrokerServer.scala
index 5c8b16fe2ad..cd061c2c8ab 100644
--- a/core/src/main/scala/kafka/server/BrokerServer.scala
+++ b/core/src/main/scala/kafka/server/BrokerServer.scala
@@ -41,7 +41,7 @@ import 
org.apache.kafka.coordinator.share.metrics.{ShareCoordinatorMetrics, Shar
 import org.apache.kafka.coordinator.share.{ShareCoordinator, 
ShareCoordinatorRecordSerde, ShareCoordinatorService}
 import org.apache.kafka.coordinator.transaction.ProducerIdManager
 import org.apache.kafka.image.publisher.{BrokerRegistrationTracker, 
MetadataPublisher}
-import org.apache.kafka.metadata.{BrokerState, ListenerInfo, 
KRaftMetadataCache, MetadataCache, MetadataVersionConfigValidator}
+import org.apache.kafka.metadata.{BrokerState, KRaftMetadataCache, 
ListenerInfo, MetadataCache, MetadataVersionConfigValidator}
 import org.apache.kafka.metadata.publisher.{AclPublisher, 
DelegationTokenPublisher, DynamicClientQuotaPublisher, 
DynamicTopicClusterQuotaPublisher, ScramPublisher}
 import org.apache.kafka.security.{CredentialProvider, DelegationTokenManager}
 import org.apache.kafka.server.authorizer.Authorizer
@@ -55,12 +55,10 @@ import 
org.apache.kafka.server.share.persister.{DefaultStatePersister, NoOpState
 import org.apache.kafka.server.share.session.ShareSessionCache
 import org.apache.kafka.server.util.timer.{SystemTimer, SystemTimerReaper}
 import org.apache.kafka.server.util.{Deadline, FutureUtils, KafkaScheduler}
-import org.apache.kafka.server.{AssignmentsManager, BrokerFeatures, 
ClientMetricsManager, DefaultApiVersionManager, DelayedActionQueue, ProcessRole}
+import org.apache.kafka.server.{AssignmentsManager, BrokerFeatures, 
ClientMetricsManager, DefaultApiVersionManager, DelayedActionQueue, 
KRaftTopicCreator, NodeToControllerChannelManagerImpl, ProcessRole, 
RaftControllerNodeProvider}
 import org.apache.kafka.server.transaction.AddPartitionsToTxnManager
 import org.apache.kafka.storage.internals.log.LogDirFailureChannel
 import org.apache.kafka.storage.log.metrics.BrokerTopicStats
-import org.apache.kafka.server.NodeToControllerChannelManagerImpl
-import org.apache.kafka.server.RaftControllerNodeProvider
 
 import java.time.Duration
 import java.util
diff --git a/core/src/main/scala/kafka/server/ForwardingManager.scala 
b/core/src/main/scala/kafka/server/ForwardingManager.scala
index 7737d2d2171..08a2d9d3ede 100644
--- a/core/src/main/scala/kafka/server/ForwardingManager.scala
+++ b/core/src/main/scala/kafka/server/ForwardingManager.scala
@@ -24,13 +24,13 @@ import org.apache.kafka.clients.{ClientResponse, 
NodeApiVersions}
 import org.apache.kafka.common.errors.TimeoutException
 import org.apache.kafka.common.metrics.Metrics
 import org.apache.kafka.common.protocol.Errors
-import org.apache.kafka.common.requests.{AbstractRequest, AbstractResponse, 
EnvelopeRequest, EnvelopeResponse, RequestContext, RequestHeader}
+import org.apache.kafka.common.requests.{AbstractRequest, AbstractResponse, 
EnvelopeResponse, RequestContext, RequestHeader}
+import org.apache.kafka.server.ForwardingManagerUtil
 import org.apache.kafka.server.common.{ControllerRequestCompletionHandler, 
NodeToControllerChannelManager}
 import org.apache.kafka.server.metrics.ForwardingManagerMetrics
 
 import java.util.Optional
 import java.util.concurrent.TimeUnit
-import scala.jdk.OptionConverters.RichOptional
 
 trait ForwardingManager {
   def close(): Unit
@@ -90,29 +90,6 @@ trait ForwardingManager {
   def controllerApiVersions: Optional[NodeApiVersions]
 }
 
-object ForwardingManager {
-  def apply(
-    channelManager: NodeToControllerChannelManager,
-    metrics: Metrics
-  ): ForwardingManager = {
-    new ForwardingManagerImpl(channelManager, metrics)
-  }
-
-  private[server] def buildEnvelopeRequest(context: RequestContext,
-                                           forwardRequestBuffer: ByteBuffer): 
EnvelopeRequest.Builder = {
-    val principalSerde = context.principalSerde.toScala.getOrElse(
-      throw new IllegalArgumentException(s"Cannot deserialize principal from 
request context $context " +
-        "since there is no serde defined")
-    )
-    val serializedPrincipal = principalSerde.serialize(context.principal)
-    new EnvelopeRequest.Builder(
-      forwardRequestBuffer,
-      serializedPrincipal,
-      context.clientAddress.getAddress
-    )
-  }
-}
-
 class ForwardingManagerImpl(
   channelManager: NodeToControllerChannelManager,
   metrics: Metrics
@@ -128,7 +105,7 @@ class ForwardingManagerImpl(
     requestToString: () => String,
     responseCallback: Option[AbstractResponse] => Unit
   ): Unit = {
-    val envelopeRequest = 
ForwardingManager.buildEnvelopeRequest(requestContext, requestBufferCopy)
+    val envelopeRequest = 
ForwardingManagerUtil.buildEnvelopeRequest(requestContext, requestBufferCopy)
     val requestCreationTimeMs = 
TimeUnit.NANOSECONDS.toMillis(requestCreationNs)
 
     class ForwardingResponseHandler extends ControllerRequestCompletionHandler 
{
diff --git 
a/core/src/test/scala/unit/kafka/server/AutoTopicCreationManagerTest.scala 
b/core/src/test/scala/unit/kafka/server/AutoTopicCreationManagerTest.scala
index 1266176f63a..5deaa94306b 100644
--- a/core/src/test/scala/unit/kafka/server/AutoTopicCreationManagerTest.scala
+++ b/core/src/test/scala/unit/kafka/server/AutoTopicCreationManagerTest.scala
@@ -40,6 +40,7 @@ import org.apache.kafka.coordinator.share.{ShareCoordinator, 
ShareCoordinatorCon
 import org.apache.kafka.metadata.MetadataCache
 import org.apache.kafka.server.config.ServerConfigs
 import org.apache.kafka.coordinator.transaction.TransactionLogConfig
+import org.apache.kafka.server.TopicCreator
 import org.apache.kafka.server.quota.ControllerMutationQuota
 import org.junit.jupiter.api.Assertions.{assertEquals, assertTrue}
 import org.junit.jupiter.api.{BeforeEach, Test}
diff --git 
a/server/src/main/java/org/apache/kafka/server/ForwardingManagerUtil.java 
b/server/src/main/java/org/apache/kafka/server/ForwardingManagerUtil.java
new file mode 100644
index 00000000000..c7281c5f696
--- /dev/null
+++ b/server/src/main/java/org/apache/kafka/server/ForwardingManagerUtil.java
@@ -0,0 +1,41 @@
+/*
+ * 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.kafka.server;
+
+import org.apache.kafka.common.requests.EnvelopeRequest;
+import org.apache.kafka.common.requests.RequestContext;
+import org.apache.kafka.common.security.auth.KafkaPrincipalSerde;
+
+import java.nio.ByteBuffer;
+
+public final class ForwardingManagerUtil {
+    public static EnvelopeRequest.Builder buildEnvelopeRequest(RequestContext 
context, ByteBuffer forwardRequestBuffer) {
+        KafkaPrincipalSerde principalSerde = 
context.principalSerde.orElseThrow(() ->
+            new IllegalArgumentException(
+                "Cannot deserialize principal from request context " + context 
+
+                " since there is no serde defined"
+            )
+        );
+        byte[] serializedPrincipal = 
principalSerde.serialize(context.principal);
+        return new EnvelopeRequest.Builder(
+            forwardRequestBuffer,
+            serializedPrincipal,
+            context.clientAddress.getAddress()
+        );
+    }
+}
diff --git a/core/src/main/java/kafka/server/KRaftTopicCreator.java 
b/server/src/main/java/org/apache/kafka/server/KRaftTopicCreator.java
similarity index 98%
rename from core/src/main/java/kafka/server/KRaftTopicCreator.java
rename to server/src/main/java/org/apache/kafka/server/KRaftTopicCreator.java
index 6cb0c63cbcb..7e730e07f80 100644
--- a/core/src/main/java/kafka/server/KRaftTopicCreator.java
+++ b/server/src/main/java/org/apache/kafka/server/KRaftTopicCreator.java
@@ -15,7 +15,7 @@
  * limitations under the License.
  */
 
-package kafka.server;
+package org.apache.kafka.server;
 
 import org.apache.kafka.clients.ClientResponse;
 import org.apache.kafka.common.errors.TimeoutException;
@@ -65,7 +65,7 @@ public class KRaftTopicCreator implements TopicCreator {
             requestContext.correlationId()
         );
 
-        AbstractRequest.Builder<? extends AbstractRequest> envelopeRequest = 
ForwardingManager$.MODULE$.buildEnvelopeRequest(
+        AbstractRequest.Builder<? extends AbstractRequest> envelopeRequest = 
ForwardingManagerUtil.buildEnvelopeRequest(
             requestContext,
             createTopicsRequest.build(requestHeader.apiVersion())
                 .serializeWithHeader(requestHeader)
diff --git a/core/src/main/java/kafka/server/TopicCreator.java 
b/server/src/main/java/org/apache/kafka/server/TopicCreator.java
similarity index 98%
rename from core/src/main/java/kafka/server/TopicCreator.java
rename to server/src/main/java/org/apache/kafka/server/TopicCreator.java
index 236de0f0900..077d8de0002 100644
--- a/core/src/main/java/kafka/server/TopicCreator.java
+++ b/server/src/main/java/org/apache/kafka/server/TopicCreator.java
@@ -15,7 +15,7 @@
  * limitations under the License.
  */
 
-package kafka.server;
+package org.apache.kafka.server;
 
 import org.apache.kafka.common.requests.CreateTopicsRequest;
 import org.apache.kafka.common.requests.CreateTopicsResponse;
diff --git a/core/src/test/java/kafka/server/KRaftTopicCreatorTest.java 
b/server/src/test/java/org/apache/kafka/server/KRaftTopicCreatorTest.java
similarity index 99%
rename from core/src/test/java/kafka/server/KRaftTopicCreatorTest.java
rename to 
server/src/test/java/org/apache/kafka/server/KRaftTopicCreatorTest.java
index 0a749bde44e..24c520d6ae7 100644
--- a/core/src/test/java/kafka/server/KRaftTopicCreatorTest.java
+++ b/server/src/test/java/org/apache/kafka/server/KRaftTopicCreatorTest.java
@@ -15,7 +15,7 @@
  * limitations under the License.
  */
 
-package kafka.server;
+package org.apache.kafka.server;
 
 import org.apache.kafka.clients.ClientResponse;
 import org.apache.kafka.clients.NodeApiVersions;

Reply via email to