@@ -257,13 +257,16 @@ internal class BigDataContainerFactory(
257257
258258 private fun kafka (): BigDataServiceContainer {
259259 val kafka = options.kafka
260- if (! kafka.kerberos.enabled) return plaintextKafka(kafka)
260+ if (! kafka.kerberos.enabled) {
261+ return if (kafka.tls.enabled) tlsKafka(kafka) else plaintextKafka(kafka)
262+ }
261263
262264 val kafkaHostPort = options.portBindings.hostPort(9092 , options.portBindings.kafka)
265+ val externalProtocol = if (kafka.tls.enabled) " SASL_SSL" else " SASL_PLAINTEXT"
263266 val advertisedListener = if (kafkaHostPort > 0 ) {
264- " SASL_PLAINTEXT ://localhost:$kafkaHostPort "
267+ " $externalProtocol ://localhost:$kafkaHostPort "
265268 } else {
266- " SASL_PLAINTEXT ://broker1.example.com:9092"
269+ " $externalProtocol ://broker1.example.com:9092"
267270 }
268271 val container = GenericBigDataContainer (kafka.image)
269272 .withNetwork(network)
@@ -284,11 +287,14 @@ internal class BigDataContainerFactory(
284287 .waitingFor(Wait .forListeningPort().withStartupTimeout(Duration .ofMinutes(3 )))
285288 mountKerberos(container)
286289 container
287- .withEnv(" KAFKA_LISTENER_SECURITY_PROTOCOL_MAP" , " CONTROLLER:PLAINTEXT,SASL_PLAINTEXT:SASL_PLAINTEXT,PLAINTEXT:PLAINTEXT" )
290+ .withEnv(
291+ " KAFKA_LISTENER_SECURITY_PROTOCOL_MAP" ,
292+ " CONTROLLER:PLAINTEXT,$externalProtocol :$externalProtocol ,PLAINTEXT:PLAINTEXT" ,
293+ )
288294 .withEnv(" KAFKA_ADVERTISED_LISTENERS" , " $advertisedListener ,PLAINTEXT://kafka:19092" )
289- .withEnv(" KAFKA_LISTENERS" , " SASL_PLAINTEXT ://0.0.0.0:9092,PLAINTEXT://0.0.0.0:19092,CONTROLLER://kafka:29093" )
295+ .withEnv(" KAFKA_LISTENERS" , " $externalProtocol ://0.0.0.0:9092,PLAINTEXT://0.0.0.0:19092,CONTROLLER://kafka:29093" )
290296 .withEnv(" KAFKA_CONTROLLER_QUORUM_VOTERS" , " 1@kafka:29093" )
291- .withEnv(" KAFKA_INTER_BROKER_LISTENER_NAME" , " SASL_PLAINTEXT " )
297+ .withEnv(" KAFKA_INTER_BROKER_LISTENER_NAME" , externalProtocol )
292298 .withEnv(" KAFKA_SASL_ENABLED_MECHANISMS" , " GSSAPI" )
293299 .withEnv(" KAFKA_SASL_MECHANISM_INTER_BROKER_PROTOCOL" , " GSSAPI" )
294300 .withEnv(" KAFKA_SASL_KERBEROS_SERVICE_NAME" , kafka.kerberos.servicePrincipal.substringBefore(" /" ))
@@ -298,6 +304,7 @@ internal class BigDataContainerFactory(
298304 Transferable .of(kafkaJaas(kafka.kerberos)),
299305 " /etc/kafka/kerberos/kafka_server_jaas.conf" ,
300306 )
307+ val sslProperties = if (kafka.tls.enabled) configureKafkaBrokerTls(container, kafka, " SASL_SSL" ) else emptyMap()
301308
302309 return BigDataServiceContainer (BigDataService .KAFKA , attachLogs(" kafka" , container)) {
303310 val bootstrapServers = " ${container.host} :${container.getMappedPort(9092 )} "
@@ -308,7 +315,41 @@ internal class BigDataContainerFactory(
308315 properties = mapOf (
309316 " bootstrap.servers" to bootstrapServers,
310317 " spring.kafka.bootstrap-servers" to bootstrapServers,
311- ) + kerberosProperties(" kafka" , kafka.kerberos) + kafkaClientKerberosProperties(kafka.kerberos, options.kerberos),
318+ ) + kerberosProperties(" kafka" , kafka.kerberos) +
319+ kafkaClientKerberosProperties(kafka.kerberos, options.kerberos, kafka.tls.enabled) +
320+ sslProperties,
321+ )
322+ }
323+ }
324+
325+ private fun tlsKafka (kafka : KafkaOptions ): BigDataServiceContainer {
326+ val kafkaHostPort = options.portBindings.hostPort(9092 , options.portBindings.kafka)
327+ val container = if (kafkaHostPort == 0 ) {
328+ KafkaContainer (DockerImageName .parse(kafka.image))
329+ } else {
330+ FixedPortKafkaContainer (DockerImageName .parse(kafka.image)).withServicePort(9092 , kafkaHostPort)
331+ }
332+ container
333+ .withNetwork(network)
334+ .withNetworkAliases(" kafka" )
335+ .withStartupTimeout(Duration .ofMinutes(3 ))
336+ .withEnv(" KAFKA_LISTENER_SECURITY_PROTOCOL_MAP" , " BROKER:SSL,PLAINTEXT:SSL,CONTROLLER:PLAINTEXT" )
337+ .withEnv(" KAFKA_INTER_BROKER_LISTENER_NAME" , " BROKER" )
338+ .withEnv(" KAFKA_SSL_CLIENT_AUTH" , " none" )
339+ .withEnv(" KAFKA_SSL_ENDPOINT_IDENTIFICATION_ALGORITHM" , " " )
340+ val sslProperties = configureKafkaBrokerTls(container, kafka, " SSL" )
341+
342+ return BigDataServiceContainer (BigDataService .KAFKA , attachLogs(" kafka" , container)) {
343+ val bootstrapServers = container.bootstrapServers
344+ BigDataEndpoint (
345+ service = BigDataService .KAFKA ,
346+ host = container.host,
347+ ports = mapOf (" bootstrap" to container.getMappedPort(9092 )),
348+ properties = mapOf (
349+ " bootstrap.servers" to bootstrapServers,
350+ " spring.kafka.bootstrap-servers" to bootstrapServers,
351+ " bootstrap.servers.internal" to " kafka:9093" ,
352+ ) + sslProperties,
312353 )
313354 }
314355 }
@@ -351,6 +392,16 @@ internal class BigDataContainerFactory(
351392 .withEnv(" SCHEMA_REGISTRY_LISTENERS" , " http://0.0.0.0:8085" )
352393 .withEnv(" SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS" , " PLAINTEXT://kafka:19092" )
353394 .waitingFor(Wait .forHttp(" /subjects" ).forStatusCode(200 ).withStartupTimeout(Duration .ofMinutes(3 )))
395+ if (kafka.tls.enabled && ! kafka.kerberos.enabled) {
396+ container
397+ .withFileSystemBind(tlsMaterial.trustStorePath.toString(), " /etc/schema-registry/tls/truststore.p12" , BindMode .READ_ONLY )
398+ .withEnv(" SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS" , " SSL://kafka:9093" )
399+ .withEnv(" SCHEMA_REGISTRY_KAFKASTORE_SECURITY_PROTOCOL" , " SSL" )
400+ .withEnv(" SCHEMA_REGISTRY_KAFKASTORE_SSL_TRUSTSTORE_LOCATION" , " /etc/schema-registry/tls/truststore.p12" )
401+ .withEnv(" SCHEMA_REGISTRY_KAFKASTORE_SSL_TRUSTSTORE_PASSWORD" , tlsMaterial.trustStorePassword)
402+ .withEnv(" SCHEMA_REGISTRY_KAFKASTORE_SSL_TRUSTSTORE_TYPE" , " PKCS12" )
403+ .withEnv(" SCHEMA_REGISTRY_KAFKASTORE_SSL_ENDPOINT_IDENTIFICATION_ALGORITHM" , " " )
404+ }
354405 return BigDataServiceContainer (BigDataService .SCHEMA_REGISTRY , attachLogs(" schema-registry" , container)) {
355406 val tlsEndpoint = httpTlsEndpoint(
356407 name = " schema-registry" ,
@@ -384,11 +435,22 @@ internal class BigDataContainerFactory(
384435 container
385436 .withEnv(" JAVA_OPTS" , " -Djava.security.krb5.conf=/kerby/client/krb5.conf" )
386437 .withEnv(" KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS" , " broker1.example.com:9092" )
387- .withEnv(" KAFKA_CLUSTERS_0_PROPERTIES_SECURITY_PROTOCOL" , " SASL_PLAINTEXT" )
438+ .withEnv(" KAFKA_CLUSTERS_0_PROPERTIES_SECURITY_PROTOCOL" , if (kafka.tls.enabled) " SASL_SSL " else " SASL_PLAINTEXT" )
388439 .withEnv(" KAFKA_CLUSTERS_0_PROPERTIES_SASL_MECHANISM" , " GSSAPI" )
389440 .withEnv(" KAFKA_CLUSTERS_0_PROPERTIES_SASL_KERBEROS_SERVICE_NAME" , " kafka" )
390441 .withEnv(" KAFKA_CLUSTERS_0_PROPERTIES_SASL_JAAS_CONFIG" , inlineJaas(kafka.kafkaUiKerberos))
391442 }
443+ if (kafka.tls.enabled) {
444+ val trustStore = tlsMaterial.trustStorePath
445+ container
446+ .withFileSystemBind(trustStore.toString(), " /etc/kafka/tls/truststore.p12" , BindMode .READ_ONLY )
447+ .withEnv(" KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS" , if (kafka.kerberos.enabled) " broker1.example.com:9092" else " kafka:9093" )
448+ .withEnv(" KAFKA_CLUSTERS_0_PROPERTIES_SECURITY_PROTOCOL" , if (kafka.kerberos.enabled) " SASL_SSL" else " SSL" )
449+ .withEnv(" KAFKA_CLUSTERS_0_PROPERTIES_SSL_TRUSTSTORE_LOCATION" , " /etc/kafka/tls/truststore.p12" )
450+ .withEnv(" KAFKA_CLUSTERS_0_PROPERTIES_SSL_TRUSTSTORE_PASSWORD" , tlsMaterial.trustStorePassword)
451+ .withEnv(" KAFKA_CLUSTERS_0_PROPERTIES_SSL_TRUSTSTORE_TYPE" , " PKCS12" )
452+ .withEnv(" KAFKA_CLUSTERS_0_PROPERTIES_SSL_ENDPOINT_IDENTIFICATION_ALGORITHM" , " " )
453+ }
392454
393455 return BigDataServiceContainer (BigDataService .KAFKA_UI , attachLogs(" kafka-ui" , container)) {
394456 val tlsEndpoint = httpTlsEndpoint(
@@ -629,6 +691,35 @@ internal class BigDataContainerFactory(
629691 return configuredHostPort
630692 }
631693
694+ private fun configureKafkaBrokerTls (
695+ container : GenericContainer <* >,
696+ kafka : KafkaOptions ,
697+ securityProtocol : String ,
698+ ): Map <String , String > {
699+ val keyStore = tlsMaterial.keyStore(
700+ name = " kafka" ,
701+ domain = kafka.tls.domain,
702+ sanDomains = listOf (" kafka" , " broker1.example.com" ),
703+ )
704+ container
705+ .withFileSystemBind(keyStore.path.toString(), " /etc/kafka/tls/kafka.keystore.p12" , BindMode .READ_ONLY )
706+ .withFileSystemBind(tlsMaterial.trustStorePath.toString(), " /etc/kafka/tls/kafka.truststore.p12" , BindMode .READ_ONLY )
707+ .withEnv(" KAFKA_SSL_KEYSTORE_LOCATION" , " /etc/kafka/tls/kafka.keystore.p12" )
708+ .withEnv(" KAFKA_SSL_KEYSTORE_PASSWORD" , keyStore.password)
709+ .withEnv(" KAFKA_SSL_KEYSTORE_TYPE" , keyStore.type)
710+ .withEnv(" KAFKA_SSL_KEY_PASSWORD" , keyStore.password)
711+ .withEnv(" KAFKA_SSL_TRUSTSTORE_LOCATION" , " /etc/kafka/tls/kafka.truststore.p12" )
712+ .withEnv(" KAFKA_SSL_TRUSTSTORE_PASSWORD" , tlsMaterial.trustStorePassword)
713+ .withEnv(" KAFKA_SSL_TRUSTSTORE_TYPE" , " PKCS12" )
714+
715+ return mapOf (
716+ " security.protocol" to securityProtocol,
717+ " ssl.truststore.location" to tlsMaterial.trustStorePath.toString(),
718+ " ssl.truststore.password" to tlsMaterial.trustStorePassword,
719+ " ssl.truststore.type" to " PKCS12" ,
720+ ) + tlsMaterial.properties()
721+ }
722+
632723 private data class HttpTlsEndpoint (
633724 val host : String? = null ,
634725 val port : Int? = null ,
@@ -758,11 +849,15 @@ internal class BigDataContainerFactory(
758849 emptyMap()
759850 }
760851
761- private fun kafkaClientKerberosProperties (service : KerberosAuthOptions , client : KerberosOptions ): Map <String , String > =
852+ private fun kafkaClientKerberosProperties (
853+ service : KerberosAuthOptions ,
854+ client : KerberosOptions ,
855+ tlsEnabled : Boolean = false,
856+ ): Map <String , String > =
762857 if (service.enabled) {
763858 val clientKeytab = localKerberosPath(" /kerby/keytabs/client.keytab" )
764859 mapOf (
765- " security.protocol" to " SASL_PLAINTEXT" ,
860+ " security.protocol" to if (tlsEnabled) " SASL_SSL " else " SASL_PLAINTEXT" ,
766861 " sasl.mechanism" to " GSSAPI" ,
767862 " sasl.kerberos.service.name" to service.servicePrincipal.substringBefore(" /" ),
768863 " sasl.jaas.config" to inlineJaas(client.clientPrincipal, clientKeytab),
0 commit comments