@@ -57,25 +57,20 @@ use crate::{
5757 node_id_hasher:: node_id_hash32_offset,
5858 } ,
5959 crd:: {
60- self , BROKER_ID_POD_MAP_DIR , KAFKA_HEAP_OPTS , LISTENER_BOOTSTRAP_VOLUME_NAME ,
60+ BROKER_ID_POD_MAP_DIR , KAFKA_HEAP_OPTS , LISTENER_BOOTSTRAP_VOLUME_NAME ,
6161 LISTENER_BROKER_VOLUME_NAME , LOG_DIRS_VOLUME_NAME , METRICS_PORT , METRICS_PORT_NAME ,
62- MetadataManager , STACKABLE_CONFIG_DIR , STACKABLE_DATA_DIR ,
63- STACKABLE_LISTENER_BOOTSTRAP_DIR , STACKABLE_LISTENER_BROKER_DIR , STACKABLE_LOG_CONFIG_DIR ,
64- STACKABLE_LOG_DIR ,
62+ STACKABLE_CONFIG_DIR , STACKABLE_DATA_DIR , STACKABLE_LISTENER_BOOTSTRAP_DIR ,
63+ STACKABLE_LISTENER_BROKER_DIR , STACKABLE_LOG_CONFIG_DIR , STACKABLE_LOG_DIR ,
6564 role:: {
6665 AnyConfig , KAFKA_NODE_ID_OFFSET , KafkaRole , broker:: BrokerContainer ,
6766 controller:: ControllerContainer ,
6867 } ,
6968 security:: KafkaTlsSecurity ,
70- v1alpha1,
7169 } ,
7270} ;
7371
7472#[ derive( Snafu , Debug ) ]
7573pub enum Error {
76- #[ snafu( display( "invalid metadata manager" ) ) ]
77- InvalidMetadataManager { source : crate :: crd:: Error } ,
78-
7974 #[ snafu( display( "failed to add kerberos config" ) ) ]
8075 AddKerberosConfig {
8176 source : crate :: controller:: build:: kerberos:: Error ,
@@ -128,9 +123,6 @@ pub enum Error {
128123 #[ snafu( display( "missing secret lifetime" ) ) ]
129124 MissingSecretLifetime ,
130125
131- #[ snafu( display( "failed to retrieve rolegroup replicas" ) ) ]
132- RoleGroupReplicas { source : crd:: role:: Error } ,
133-
134126 #[ snafu( display( "vector agent is enabled but vector aggregator ConfigMap is missing" ) ) ]
135127 VectorAggregatorConfigMapMissing ,
136128}
@@ -140,7 +132,6 @@ pub enum Error {
140132/// The [`Pod`](`stackable_operator::k8s_openapi::api::core::v1::Pod`)s are accessible through the corresponding
141133/// [`Service`](`stackable_operator::k8s_openapi::api::core::v1::Service`) from [`build_rolegroup_headless_service`](`crate::controller::build::resource::service::build_rolegroup_headless_service`).
142134pub fn build_broker_rolegroup_statefulset (
143- kafka : & v1alpha1:: KafkaCluster ,
144135 kafka_role : & KafkaRole ,
145136 role_group_name : & RoleGroupName ,
146137 validated_cluster : & ValidatedCluster ,
@@ -213,7 +204,9 @@ pub fn build_broker_rolegroup_statefulset(
213204
214205 let mut env = Vec :: < EnvVar > :: from ( validated_rg. env_overrides . clone ( ) ) ;
215206
216- if let Some ( zookeeper_config_map_name) = & kafka. spec . cluster_config . zookeeper_config_map_name {
207+ if let Some ( zookeeper_config_map_name) =
208+ & validated_cluster. cluster_config . zookeeper_config_map_name
209+ {
217210 env. push ( EnvVar {
218211 name : "ZOOKEEPER" . to_string ( ) ,
219212 value_from : Some ( EnvVarSource {
@@ -240,10 +233,6 @@ pub fn build_broker_rolegroup_statefulset(
240233 ..EnvVar :: default ( )
241234 } ) ;
242235
243- let metadata_manager = kafka
244- . effective_metadata_manager ( )
245- . context ( InvalidMetadataManagerSnafu ) ?;
246-
247236 cb_kafka
248237 . image_from_product_image ( resolved_product_image)
249238 . command ( vec ! [
@@ -254,7 +243,7 @@ pub fn build_broker_rolegroup_statefulset(
254243 "-c" . to_string( ) ,
255244 ] )
256245 . args ( vec ! [ broker_kafka_container_commands(
257- metadata_manager == MetadataManager :: KRaft ,
246+ validated_cluster . cluster_config . is_kraft_mode ( ) ,
258247 // we need controller pods
259248 validated_cluster
260249 . pod_descriptors( Some ( & KafkaRole :: Controller ) )
@@ -355,8 +344,9 @@ pub fn build_broker_rolegroup_statefulset(
355344 . context ( AddListenerVolumeSnafu ) ?;
356345 }
357346
358- if let Some ( broker_id_config_map_name) =
359- & kafka. spec . cluster_config . broker_id_pod_config_map_name
347+ if let Some ( broker_id_config_map_name) = & validated_cluster
348+ . cluster_config
349+ . broker_id_pod_config_map_name
360350 {
361351 pod_builder
362352 . add_volume (
@@ -381,7 +371,7 @@ pub fn build_broker_rolegroup_statefulset(
381371
382372 add_vector_container (
383373 & mut pod_builder,
384- kafka ,
374+ validated_cluster ,
385375 resolved_product_image,
386376 merged_config,
387377 ) ?;
@@ -411,10 +401,7 @@ pub fn build_broker_rolegroup_statefulset(
411401 . build ( ) ,
412402 spec : Some ( StatefulSetSpec {
413403 pod_management_policy : Some ( "Parallel" . to_string ( ) ) ,
414- replicas : kafka_role
415- . replicas ( kafka, role_group_name. as_ref ( ) )
416- . context ( RoleGroupReplicasSnafu ) ?
417- . map ( i32:: from) ,
404+ replicas : Some ( i32:: from ( validated_rg. replicas ) ) ,
418405 selector : LabelSelector {
419406 match_labels : Some (
420407 validated_cluster
@@ -434,7 +421,6 @@ pub fn build_broker_rolegroup_statefulset(
434421
435422/// The controller rolegroup [`StatefulSet`] runs the rolegroup, as configured by the administrator.
436423pub fn build_controller_rolegroup_statefulset (
437- kafka : & v1alpha1:: KafkaCluster ,
438424 kafka_role : & KafkaRole ,
439425 role_group_name : & RoleGroupName ,
440426 validated_cluster : & ValidatedCluster ,
@@ -500,7 +486,9 @@ pub fn build_controller_rolegroup_statefulset(
500486 } ) ;
501487
502488 // Controllers need the ZooKeeper connection string for migration
503- if let Some ( zookeeper_config_map_name) = & kafka. spec . cluster_config . zookeeper_config_map_name {
489+ if let Some ( zookeeper_config_map_name) =
490+ & validated_cluster. cluster_config . zookeeper_config_map_name
491+ {
504492 env. push ( EnvVar {
505493 name : "ZOOKEEPER" . to_string ( ) ,
506494 value_from : Some ( EnvVarSource {
@@ -606,7 +594,7 @@ pub fn build_controller_rolegroup_statefulset(
606594
607595 add_vector_container (
608596 & mut pod_builder,
609- kafka ,
597+ validated_cluster ,
610598 resolved_product_image,
611599 merged_config,
612600 ) ?;
@@ -636,10 +624,7 @@ pub fn build_controller_rolegroup_statefulset(
636624 type_ : Some ( "RollingUpdate" . to_string ( ) ) ,
637625 ..StatefulSetUpdateStrategy :: default ( )
638626 } ) ,
639- replicas : kafka_role
640- . replicas ( kafka, role_group_name. as_ref ( ) )
641- . context ( RoleGroupReplicasSnafu ) ?
642- . map ( i32:: from) ,
627+ replicas : Some ( i32:: from ( validated_rg. replicas ) ) ,
643628 selector : LabelSelector {
644629 match_labels : Some (
645630 validated_cluster
@@ -795,13 +780,16 @@ fn add_common_pod_config(
795780/// configured on the cluster.
796781fn add_vector_container (
797782 pod_builder : & mut PodBuilder ,
798- kafka : & v1alpha1 :: KafkaCluster ,
783+ validated_cluster : & ValidatedCluster ,
799784 resolved_product_image : & ResolvedProductImage ,
800785 merged_config : & AnyConfig ,
801786) -> Result < ( ) , Error > {
802787 // Add vector container after kafka container to keep the defaulting into kafka container
803788 if merged_config. vector_logging_enabled ( ) {
804- match & kafka. spec . cluster_config . vector_aggregator_config_map_name {
789+ match & validated_cluster
790+ . cluster_config
791+ . vector_aggregator_config_map_name
792+ {
805793 Some ( vector_aggregator_config_map_name) => {
806794 pod_builder. add_container (
807795 product_logging:: framework:: vector_container (
0 commit comments