@@ -18,8 +18,8 @@ use crate::{
1818 listener:: { KafkaListenerConfig , KafkaListenerName } ,
1919 role:: {
2020 AnyConfig , KAFKA_ADVERTISED_LISTENERS , KAFKA_CONTROLLER_QUORUM_BOOTSTRAP_SERVERS ,
21- KAFKA_LISTENER_SECURITY_PROTOCOL_MAP , KAFKA_LISTENERS , KAFKA_LOG_DIRS , KAFKA_NODE_ID ,
22- KAFKA_PROCESS_ROLES , KafkaRole ,
21+ KAFKA_CONTROLLER_QUORUM_VOTERS , KAFKA_LISTENER_SECURITY_PROTOCOL_MAP , KAFKA_LISTENERS ,
22+ KAFKA_LOG_DIRS , KAFKA_NODE_ID , KAFKA_PROCESS_ROLES , KafkaRole ,
2323 } ,
2424 security:: KafkaTlsSecurity ,
2525 v1alpha1,
@@ -94,6 +94,7 @@ pub fn build_rolegroup_config_map(
9494 pod_descriptors,
9595 listener_config,
9696 opa_connect_string,
97+ resolved_product_image. product_version . starts_with ( "3.7" ) , // needs_quorum_voters
9798 ) ?;
9899
99100 // Need to call this to get configOverrides :(
@@ -198,6 +199,7 @@ fn server_properties_file(
198199 pod_descriptors : & [ KafkaPodDescriptor ] ,
199200 listener_config : & KafkaListenerConfig ,
200201 opa_connect_string : Option < & str > ,
202+ needs_quorum_voters : bool ,
201203) -> Result < BTreeMap < String , String > , Error > {
202204 let kraft_controllers = kraft_controllers ( pod_descriptors) ;
203205
@@ -209,7 +211,7 @@ fn server_properties_file(
209211 KafkaRole :: Controller => {
210212 let kraft_controllers = kraft_controllers. context ( NoKraftControllersFoundSnafu ) ?;
211213
212- Ok ( BTreeMap :: from ( [
214+ let mut result = BTreeMap :: from ( [
213215 (
214216 KAFKA_LOG_DIRS . to_string ( ) ,
215217 "/stackable/data/kraft" . to_string ( ) ,
@@ -227,10 +229,6 @@ fn server_properties_file(
227229 KAFKA_CONTROLLER_QUORUM_BOOTSTRAP_SERVERS . to_string ( ) ,
228230 kraft_controllers. clone ( ) ,
229231 ) ,
230- // TODO: figure this out
231- //(KAFKA_CONTROLLER_QUORUM_VOTERS.to_string(),
232- //kraft_controllers,
233- //),
234232 (
235233 KAFKA_LISTENERS . to_string ( ) ,
236234 "CONTROLLER://${env:POD_NAME}.${env:ROLEGROUP_HEADLESS_SERVICE_NAME}.${env:NAMESPACE}.svc.${env:CLUSTER_DOMAIN}:${env:KAFKA_CLIENT_PORT}" . to_string ( ) ,
@@ -240,7 +238,16 @@ fn server_properties_file(
240238 listener_config
241239 . listener_security_protocol_map_for_listener ( & KafkaListenerName :: Controller )
242240 . unwrap_or ( "" . to_string ( ) ) ) ,
243- ] ) )
241+ ] ) ;
242+
243+ if needs_quorum_voters {
244+ let kraft_voters =
245+ kraft_voters ( pod_descriptors) . context ( NoKraftControllersFoundSnafu ) ?;
246+
247+ result. extend ( [ ( KAFKA_CONTROLLER_QUORUM_VOTERS . to_string ( ) , kraft_voters) ] ) ;
248+ }
249+
250+ Ok ( result)
244251 }
245252 KafkaRole :: Broker => {
246253 let mut result = BTreeMap :: from ( [
@@ -278,6 +285,13 @@ fn server_properties_file(
278285 kraft_controllers. clone ( ) ,
279286 ) ,
280287 ] ) ;
288+
289+ if needs_quorum_voters {
290+ let kraft_voters =
291+ kraft_voters ( pod_descriptors) . context ( NoKraftControllersFoundSnafu ) ?;
292+
293+ result. extend ( [ ( KAFKA_CONTROLLER_QUORUM_VOTERS . to_string ( ) , kraft_voters) ] ) ;
294+ }
281295 } else {
282296 // Running with ZooKeeper enabled
283297 result. extend ( [ (
@@ -313,7 +327,28 @@ fn kraft_controllers(pod_descriptors: &[KafkaPodDescriptor]) -> Option<String> {
313327 let result = pod_descriptors
314328 . iter ( )
315329 . filter ( |pd| pd. role == KafkaRole :: Controller . to_string ( ) )
316- . map ( |desc| format ! ( "{fqdn}:${{env:KAFKA_CLIENT_PORT}}" , fqdn = desc. fqdn( ) ) )
330+ . map ( |desc| {
331+ format ! (
332+ "{fqdn}:{client_port}" ,
333+ fqdn = desc. fqdn( ) ,
334+ client_port = desc. client_port
335+ )
336+ } )
337+ . collect :: < Vec < String > > ( )
338+ . join ( "," ) ;
339+
340+ if result. is_empty ( ) {
341+ None
342+ } else {
343+ Some ( result)
344+ }
345+ }
346+
347+ fn kraft_voters ( pod_descriptors : & [ KafkaPodDescriptor ] ) -> Option < String > {
348+ let result = pod_descriptors
349+ . iter ( )
350+ . filter ( |pd| pd. role == KafkaRole :: Controller . to_string ( ) )
351+ . map ( |desc| desc. as_quorum_voter ( ) )
317352 . collect :: < Vec < String > > ( )
318353 . join ( "," ) ;
319354
0 commit comments