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;