11//! Ensures that `Pod`s are configured and running for each [`v1alpha1::KafkaCluster`].
22
3- use std:: sync:: Arc ;
3+ use std:: { collections :: BTreeMap , sync:: Arc } ;
44
55use const_format:: concatcp;
66use snafu:: { ResultExt , Snafu } ;
77use stackable_operator:: {
88 cli:: OperatorEnvironmentOptions ,
99 cluster_resources:: { ClusterResourceApplyStrategy , ClusterResources } ,
10- commons:: rbac:: build_rbac_resources,
10+ commons:: { product_image_selection :: ResolvedProductImage , rbac:: build_rbac_resources} ,
1111 crd:: listener,
1212 kube:: {
1313 Resource ,
@@ -32,8 +32,10 @@ mod validate;
3232use crate :: {
3333 crd:: {
3434 self , APP_NAME , KafkaClusterStatus , OPERATOR_NAME ,
35+ authorization:: KafkaAuthorizationConfig ,
3536 listener:: get_kafka_listener_config,
3637 role:: { AnyConfig , KafkaRole } ,
38+ security:: KafkaTlsSecurity ,
3739 v1alpha1,
3840 } ,
3941 discovery:: { self , build_discovery_configmap} ,
@@ -208,6 +210,35 @@ impl ReconcilerError for Error {
208210 }
209211}
210212
213+ /// The validated cluster. Carries everything the build steps need, resolved once
214+ /// here so downstream code never re-derives it or touches the raw spec.
215+ pub struct ValidatedKafkaCluster {
216+ pub image : ResolvedProductImage ,
217+ pub kafka_security : KafkaTlsSecurity ,
218+ // DESIGN DECISION: the dereferenced authorization config is folded into the
219+ // validated cluster (read from here downstream). The other dereferenced input,
220+ // the authentication classes, is intentionally NOT stored: it is fully consumed
221+ // here to build `kafka_security`. Alternative: also store the resolved auth
222+ // classes — rejected because nothing downstream needs them beyond kafka_security.
223+ pub authorization_config : Option < KafkaAuthorizationConfig > ,
224+ pub role_groups : BTreeMap < KafkaRole , BTreeMap < String , ValidatedRoleGroupConfig > > ,
225+ }
226+
227+ pub struct ValidatedRoleGroupConfig {
228+ pub merged_config : AnyConfig ,
229+ // DESIGN DECISION: overrides are resolved into flat maps HERE rather than stored
230+ // as the typed KeyValueConfigOverrides and resolved in the per-file builders (the
231+ // hdfs-operator pattern). Reason: broker and controller use different override
232+ // struct types (KafkaBrokerConfigOverrides vs KafkaControllerConfigOverrides), so a
233+ // single typed field would require an enum. Resolving here keeps the build/properties
234+ // builders taking plain `BTreeMap<String,String>`. Alternative: an enum over the two
235+ // override types threaded to builders that call resolved_overrides() — more types for
236+ // no behavioural gain.
237+ pub config_file_overrides : BTreeMap < String , String > ,
238+ pub jvm_security_overrides : BTreeMap < String , String > ,
239+ pub env_overrides : BTreeMap < String , String > ,
240+ }
241+
211242pub async fn reconcile_kafka (
212243 kafka : Arc < DeserializeGuard < v1alpha1:: KafkaCluster > > ,
213244 ctx : Arc < Ctx > ,
@@ -228,15 +259,12 @@ pub async fn reconcile_kafka(
228259 . context ( DereferenceSnafu ) ?;
229260
230261 // validate (no client required)
231- let validate:: ValidatedKafkaCluster {
232- authorization_config,
233- image,
234- kafka_security,
235- role_groups,
236- } = validate:: validate ( kafka, dereferenced_objects, & ctx. operator_environment )
237- . context ( ValidateClusterSnafu ) ?;
238-
239- let opa_connect = authorization_config
262+ let validated_cluster =
263+ validate:: validate ( kafka, dereferenced_objects, & ctx. operator_environment )
264+ . context ( ValidateClusterSnafu ) ?;
265+
266+ let opa_connect = validated_cluster
267+ . authorization_config
240268 . as_ref ( )
241269 . map ( |auth_config| auth_config. opa_connect . clone ( ) ) ;
242270
@@ -251,10 +279,10 @@ pub async fn reconcile_kafka(
251279 . context ( CreateClusterResourcesSnafu ) ?;
252280
253281 tracing:: debug!(
254- kerberos_enabled = kafka_security. has_kerberos_enabled( ) ,
255- kerberos_secret_class = ?kafka_security. kerberos_secret_class( ) ,
256- tls_enabled = kafka_security. tls_enabled( ) ,
257- tls_client_authentication_class = ?kafka_security. tls_client_authentication_class( ) ,
282+ kerberos_enabled = validated_cluster . kafka_security. has_kerberos_enabled( ) ,
283+ kerberos_secret_class = ?validated_cluster . kafka_security. kerberos_secret_class( ) ,
284+ tls_enabled = validated_cluster . kafka_security. tls_enabled( ) ,
285+ tls_client_authentication_class = ?validated_cluster . kafka_security. tls_client_authentication_class( ) ,
258286 "The following security settings are used"
259287 ) ;
260288
@@ -280,20 +308,25 @@ pub async fn reconcile_kafka(
280308
281309 let mut bootstrap_listeners = Vec :: < listener:: v1alpha1:: Listener > :: new ( ) ;
282310
283- for ( kafka_role, rg_map) in & role_groups {
311+ for ( kafka_role, rg_map) in & validated_cluster . role_groups {
284312 for ( rolegroup_name, validated_rg) in rg_map {
285313 let rolegroup_ref = kafka. rolegroup_ref ( kafka_role, rolegroup_name) ;
286314
287- let rg_headless_service =
288- build_rolegroup_headless_service ( kafka, & image, & rolegroup_ref, & kafka_security)
289- . context ( BuildServiceSnafu ) ?;
315+ let rg_headless_service = build_rolegroup_headless_service (
316+ kafka,
317+ & validated_cluster. image ,
318+ & rolegroup_ref,
319+ & validated_cluster. kafka_security ,
320+ )
321+ . context ( BuildServiceSnafu ) ?;
290322
291- let rg_metrics_service = build_rolegroup_metrics_service ( kafka, & image, & rolegroup_ref)
292- . context ( BuildServiceSnafu ) ?;
323+ let rg_metrics_service =
324+ build_rolegroup_metrics_service ( kafka, & validated_cluster. image , & rolegroup_ref)
325+ . context ( BuildServiceSnafu ) ?;
293326
294327 let kafka_listeners = get_kafka_listener_config (
295328 kafka,
296- & kafka_security,
329+ & validated_cluster . kafka_security ,
297330 & rolegroup_ref,
298331 & client. kubernetes_cluster_info ,
299332 )
@@ -303,18 +336,15 @@ pub async fn reconcile_kafka(
303336 . pod_descriptors (
304337 None ,
305338 & client. kubernetes_cluster_info ,
306- kafka_security. client_port ( ) ,
339+ validated_cluster . kafka_security . client_port ( ) ,
307340 )
308341 . context ( BuildPodDescriptorsSnafu ) ?;
309342
310343 let rg_configmap = build:: config_map:: build_rolegroup_config_map (
311344 kafka,
312- & image,
313- & kafka_security,
345+ & validated_cluster,
314346 & rolegroup_ref,
315- validated_rg. config_file_overrides . clone ( ) ,
316- validated_rg. jvm_security_overrides . clone ( ) ,
317- & validated_rg. merged_config ,
347+ validated_rg,
318348 & kafka_listeners,
319349 & pod_descriptors,
320350 opa_connect. as_deref ( ) ,
@@ -325,23 +355,19 @@ pub async fn reconcile_kafka(
325355 KafkaRole :: Broker => build_broker_rolegroup_statefulset (
326356 kafka,
327357 kafka_role,
328- & image ,
358+ & validated_cluster ,
329359 & rolegroup_ref,
330- & validated_rg. env_overrides ,
331- & kafka_security,
332- & validated_rg. merged_config ,
360+ validated_rg,
333361 & rbac_sa,
334362 & client. kubernetes_cluster_info ,
335363 )
336364 . context ( BuildStatefulsetSnafu ) ?,
337365 KafkaRole :: Controller => build_controller_rolegroup_statefulset (
338366 kafka,
339367 kafka_role,
340- & image ,
368+ & validated_cluster ,
341369 & rolegroup_ref,
342- & validated_rg. env_overrides ,
343- & kafka_security,
344- & validated_rg. merged_config ,
370+ validated_rg,
345371 & rbac_sa,
346372 & client. kubernetes_cluster_info ,
347373 )
@@ -351,8 +377,7 @@ pub async fn reconcile_kafka(
351377 if let AnyConfig :: Broker ( broker_config) = & validated_rg. merged_config {
352378 let rg_bootstrap_listener = build_broker_rolegroup_bootstrap_listener (
353379 kafka,
354- & image,
355- & kafka_security,
380+ & validated_cluster,
356381 & rolegroup_ref,
357382 broker_config,
358383 )
@@ -409,7 +434,7 @@ pub async fn reconcile_kafka(
409434 }
410435
411436 let discovery_cm =
412- build_discovery_configmap ( kafka, kafka, & image , & kafka_security , & bootstrap_listeners)
437+ build_discovery_configmap ( kafka, kafka, validated_cluster , & bootstrap_listeners)
413438 . context ( BuildDiscoveryConfigSnafu ) ?;
414439
415440 cluster_resources
0 commit comments