Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -19,12 +19,11 @@ package kafka.server

import java.util.concurrent.LinkedBlockingDeque
import java.util.concurrent.atomic.AtomicReference

import kafka.common.{InterBrokerSendThread, RequestAndCompletionHandler}
import kafka.raft.RaftManager
import kafka.utils.Logging
import org.apache.kafka.clients._
import org.apache.kafka.common.Node
import org.apache.kafka.common.{Node, Reconfigurable}
import org.apache.kafka.common.metrics.Metrics
import org.apache.kafka.common.network._
import org.apache.kafka.common.protocol.Errors
Expand Down Expand Up @@ -189,6 +188,10 @@ class BrokerToControllerChannelManagerImpl(
config.saslInterBrokerHandshakeRequestEnable,
logContext
)
channelBuilder match {
case reconfigurable: Reconfigurable => config.addReconfigurable(reconfigurable)
case _ =>
Comment thread
divijvaidya marked this conversation as resolved.
}
val selector = new Selector(
NetworkReceive.UNLIMITED,
Selector.NO_IDLE_TIMEOUT_MS,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,6 @@ import java.time.Duration
import java.util
import java.util.{Collections, Properties}
import java.util.concurrent._

import javax.management.ObjectName
import com.yammer.metrics.core.MetricName
import kafka.admin.ConfigCommand
Expand All @@ -35,9 +34,9 @@ import kafka.controller.{ControllerBrokerStateInfo, ControllerChannelManager}
import kafka.log.{CleanerConfig, LogConfig}
import kafka.message.ProducerCompressionCodec
import kafka.network.{Processor, RequestChannel}
import kafka.server.QuorumTestHarness
import kafka.utils._
import kafka.utils.Implicits._
import kafka.utils.TestUtils.TestControllerRequestCompletionHandler
import kafka.zk.ConfigEntityChangeNotificationZNode
import org.apache.kafka.clients.CommonClientConfigs
import org.apache.kafka.clients.admin.AlterConfigOp.OpType
Expand All @@ -52,10 +51,12 @@ import org.apache.kafka.common.config.types.Password
import org.apache.kafka.common.config.provider.FileConfigProvider
import org.apache.kafka.common.errors.{AuthenticationException, InvalidRequestException}
import org.apache.kafka.common.internals.Topic
import org.apache.kafka.common.message.MetadataRequestData
import org.apache.kafka.common.metrics.{KafkaMetric, MetricsContext, MetricsReporter, Quota}
import org.apache.kafka.common.network.{ListenerName, Mode}
import org.apache.kafka.common.network.CertStores.{KEYSTORE_PROPS, TRUSTSTORE_PROPS}
import org.apache.kafka.common.record.TimestampType
import org.apache.kafka.common.requests.MetadataRequest
import org.apache.kafka.common.security.auth.SecurityProtocol
import org.apache.kafka.common.security.scram.ScramCredential
import org.apache.kafka.common.serialization.{StringDeserializer, StringSerializer}
Expand Down Expand Up @@ -429,6 +430,23 @@ class DynamicBrokerReconfigurationTest extends QuorumTestHarness with SaslSetup
verifyProduceConsume(producer, consumer, 10, topic)
}

def verifyBrokerToControllerCall(controller: KafkaServer): Unit = {
val nonControllerBroker = servers.find(_.config.brokerId != controller.config.brokerId).get
val brokerToControllerManager = nonControllerBroker.clientToControllerChannelManager
val completionHandler = new TestControllerRequestCompletionHandler()
brokerToControllerManager.sendRequest(new MetadataRequest.Builder(new MetadataRequestData()), completionHandler)
TestUtils.waitUntilTrue(() => {
completionHandler.completed.get() || completionHandler.timedOut.get()
}, "Timed out while waiting for broker to controller API call")
// we do not expect a timeout from broker to controller request
assertFalse(completionHandler.timedOut.get(), "broker to controller request is timeout")
assertTrue(completionHandler.actualResponse.isDefined, "No response recorded even though request is completed")
val response = completionHandler.actualResponse.get
assertNull(response.authenticationException(), s"Request failed due to authentication error ${response.authenticationException}")
assertNull(response.versionMismatch(), s"Request failed due to unsupported version error ${response.versionMismatch}")
assertFalse(response.wasDisconnected(), "Request failed because broker is not available")
}

// Produce/consume should work with old as well as new client keystore
verifySslProduceConsume(sslProperties1, "alter-truststore-1")
verifySslProduceConsume(sslProperties2, "alter-truststore-2")
Expand Down Expand Up @@ -469,6 +487,9 @@ class DynamicBrokerReconfigurationTest extends QuorumTestHarness with SaslSetup
JTestUtils.fieldValue(controllerChannelManager, classOf[ControllerChannelManager], "brokerStateInfo")
brokerStateInfo(0).networkClient.disconnect("0")
TestUtils.createTopic(zkClient, "testtopic2", numPartitions, replicationFactor = numServers, servers)

// validate that the brokerToController request works fine
verifyBrokerToControllerCall(controller)
}

@Test
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,13 +19,14 @@ package kafka.server

import java.nio.ByteBuffer
import java.util.Collections
import java.util.concurrent.atomic.{AtomicBoolean, AtomicReference}
import java.util.concurrent.atomic.AtomicReference
import kafka.utils.TestUtils
import kafka.utils.TestUtils.TestControllerRequestCompletionHandler
import org.apache.kafka.clients.{ClientResponse, ManualMetadataUpdater, Metadata, MockClient, NodeApiVersions}
import org.apache.kafka.common.Node
import org.apache.kafka.common.message.{EnvelopeResponseData, MetadataRequestData}
import org.apache.kafka.common.protocol.{ApiKeys, Errors}
import org.apache.kafka.common.requests.{AbstractRequest, EnvelopeRequest, EnvelopeResponse, MetadataRequest, MetadataResponse, RequestTestUtils}
import org.apache.kafka.common.requests.{AbstractRequest, EnvelopeRequest, EnvelopeResponse, MetadataRequest, RequestTestUtils}
import org.apache.kafka.common.security.auth.KafkaPrincipal
import org.apache.kafka.common.security.authenticator.DefaultKafkaPrincipalBuilder
import org.apache.kafka.common.utils.MockTime
Expand All @@ -51,7 +52,7 @@ class BrokerToControllerRequestThreadTest {
config, time, "", retryTimeoutMs)
testRequestThread.started = true

val completionHandler = new TestRequestCompletionHandler(None)
val completionHandler = new TestControllerRequestCompletionHandler(None)
val queueItem = BrokerToControllerQueueItem(
time.milliseconds(),
new MetadataRequest.Builder(new MetadataRequestData()),
Expand Down Expand Up @@ -89,7 +90,7 @@ class BrokerToControllerRequestThreadTest {
testRequestThread.started = true
mockClient.prepareResponse(expectedResponse)

val completionHandler = new TestRequestCompletionHandler(Some(expectedResponse))
val completionHandler = new TestControllerRequestCompletionHandler(Some(expectedResponse))
val queueItem = BrokerToControllerQueueItem(
time.milliseconds(),
new MetadataRequest.Builder(new MetadataRequestData()),
Expand Down Expand Up @@ -130,7 +131,7 @@ class BrokerToControllerRequestThreadTest {
controllerNodeProvider, config, time, "", retryTimeoutMs = Long.MaxValue)
testRequestThread.started = true

val completionHandler = new TestRequestCompletionHandler(Some(expectedResponse))
val completionHandler = new TestControllerRequestCompletionHandler(Some(expectedResponse))
val queueItem = BrokerToControllerQueueItem(
time.milliseconds(),
new MetadataRequest.Builder(new MetadataRequestData()),
Expand Down Expand Up @@ -180,7 +181,7 @@ class BrokerToControllerRequestThreadTest {
config, time, "", retryTimeoutMs = Long.MaxValue)
testRequestThread.started = true

val completionHandler = new TestRequestCompletionHandler(Some(expectedResponse))
val completionHandler = new TestControllerRequestCompletionHandler(Some(expectedResponse))
val queueItem = BrokerToControllerQueueItem(
time.milliseconds(),
new MetadataRequest.Builder(new MetadataRequestData()
Expand Down Expand Up @@ -243,7 +244,7 @@ class BrokerToControllerRequestThreadTest {
config, time, "", retryTimeoutMs = Long.MaxValue)
testRequestThread.started = true

val completionHandler = new TestRequestCompletionHandler(Some(expectedResponse))
val completionHandler = new TestControllerRequestCompletionHandler(Some(expectedResponse))
val kafkaPrincipal = new KafkaPrincipal(KafkaPrincipal.USER_TYPE, "principal", true)
val kafkaPrincipalBuilder = new DefaultKafkaPrincipalBuilder(null, null)

Expand Down Expand Up @@ -305,7 +306,7 @@ class BrokerToControllerRequestThreadTest {
config, time, "", retryTimeoutMs)
testRequestThread.started = true

val completionHandler = new TestRequestCompletionHandler()
val completionHandler = new TestControllerRequestCompletionHandler()
val queueItem = BrokerToControllerQueueItem(
time.milliseconds(),
new MetadataRequest.Builder(new MetadataRequestData()
Expand Down Expand Up @@ -419,7 +420,7 @@ class BrokerToControllerRequestThreadTest {
val testRequestThread = new BrokerToControllerRequestThread(mockClient, new ManualMetadataUpdater(), controllerNodeProvider,
config, time, "", retryTimeoutMs = Long.MaxValue)

val completionHandler = new TestRequestCompletionHandler(None)
val completionHandler = new TestControllerRequestCompletionHandler(None)
val queueItem = BrokerToControllerQueueItem(
time.milliseconds(),
new MetadataRequest.Builder(new MetadataRequestData()),
Expand All @@ -445,22 +446,4 @@ class BrokerToControllerRequestThreadTest {
fail(s"Condition failed to be met after polling $tries times")
}
}

class TestRequestCompletionHandler(
expectedResponse: Option[MetadataResponse] = None
) extends ControllerRequestCompletionHandler {
val completed: AtomicBoolean = new AtomicBoolean(false)
val timedOut: AtomicBoolean = new AtomicBoolean(false)

override def onComplete(response: ClientResponse): Unit = {
expectedResponse.foreach { expected =>
assertEquals(expected, response.responseBody())
}
completed.set(true)
}

override def onTimeout(): Unit = {
timedOut.set(true)
}
}
}
24 changes: 21 additions & 3 deletions core/src/test/scala/unit/kafka/utils/TestUtils.scala
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ import kafka.server.checkpoints.OffsetCheckpointFile
import kafka.server.metadata.{ConfigRepository, MockConfigRepository}
import kafka.utils.Implicits._
import kafka.zk._
import org.apache.kafka.clients.CommonClientConfigs
import org.apache.kafka.clients.{ClientResponse, CommonClientConfigs}
import org.apache.kafka.clients.admin.AlterConfigOp.OpType
import org.apache.kafka.clients.admin._
import org.apache.kafka.clients.consumer._
Expand All @@ -61,7 +61,7 @@ import org.apache.kafka.common.network.{ClientInformation, ListenerName, Mode}
import org.apache.kafka.common.protocol.{ApiKeys, Errors}
import org.apache.kafka.common.quota.{ClientQuotaAlteration, ClientQuotaEntity}
import org.apache.kafka.common.record._
import org.apache.kafka.common.requests.{AbstractRequest, EnvelopeRequest, RequestContext, RequestHeader}
import org.apache.kafka.common.requests.{AbstractRequest, AbstractResponse, EnvelopeRequest, RequestContext, RequestHeader}
import org.apache.kafka.common.resource.ResourcePattern
import org.apache.kafka.common.security.auth.{KafkaPrincipal, KafkaPrincipalSerde, SecurityProtocol}
import org.apache.kafka.common.serialization.{ByteArrayDeserializer, ByteArraySerializer, Deserializer, IntegerSerializer, Serializer}
Expand Down Expand Up @@ -2239,4 +2239,22 @@ object TestUtils extends Logging {
s"${unexpected.mkString("`", ",", "`")}")
}

}
class TestControllerRequestCompletionHandler(expectedResponse: Option[AbstractResponse] = None)
extends ControllerRequestCompletionHandler {
var actualResponse: Option[ClientResponse] = Option.empty
val completed: AtomicBoolean = new AtomicBoolean(false)
val timedOut: AtomicBoolean = new AtomicBoolean(false)

override def onComplete(response: ClientResponse): Unit = {
actualResponse = Some(response)
expectedResponse.foreach { expected =>
assertEquals(expected, response.responseBody())
}
completed.set(true)
}

override def onTimeout(): Unit = {
timedOut.set(true)
}
}
}