Skip to content

Commit 4eb786b

Browse files
committed
push pod_descriptors calc down to build_rolegroup_config_map
1 parent 7706e88 commit 4eb786b

4 files changed

Lines changed: 19 additions & 17 deletions

File tree

rust/operator-binary/src/controller.rs

Lines changed: 2 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ use stackable_operator::{
2222
compute_conditions, operations::ClusterOperationsConditionBuilder,
2323
statefulset::StatefulSetConditionBuilder,
2424
},
25+
utils::cluster_info::KubernetesClusterInfo,
2526
};
2627
use strum::{EnumDiscriminants, IntoStaticStr};
2728

@@ -65,9 +66,6 @@ pub enum Error {
6566
#[snafu(display("failed to validate cluster"))]
6667
ValidateCluster { source: validate::Error },
6768

68-
#[snafu(display("failed to build pod descriptors"))]
69-
BuildPodDescriptors { source: crate::crd::Error },
70-
7169
#[snafu(display("invalid kafka listeners"))]
7270
InvalidKafkaListeners {
7371
source: crate::crd::listener::KafkaListenerError,
@@ -205,7 +203,6 @@ impl ReconcilerError for Error {
205203
Error::BuildService { .. } => None,
206204
Error::BuildListener { .. } => None,
207205
Error::InvalidKafkaListeners { .. } => None,
208-
Error::BuildPodDescriptors { .. } => None,
209206
}
210207
}
211208
}
@@ -222,6 +219,7 @@ pub struct ValidatedKafkaCluster {
222219
// classes — rejected because nothing downstream needs them beyond kafka_security.
223220
pub authorization_config: Option<KafkaAuthorizationConfig>,
224221
pub role_groups: BTreeMap<KafkaRole, BTreeMap<String, ValidatedRoleGroupConfig>>,
222+
pub kubernetes_cluster_info: KubernetesClusterInfo,
225223
}
226224

227225
pub struct ValidatedRoleGroupConfig {
@@ -327,21 +325,12 @@ pub async fn reconcile_kafka(
327325
)
328326
.context(InvalidKafkaListenersSnafu)?;
329327

330-
let pod_descriptors = kafka
331-
.pod_descriptors(
332-
None,
333-
&client.kubernetes_cluster_info,
334-
validated_cluster.kafka_security.client_port(),
335-
)
336-
.context(BuildPodDescriptorsSnafu)?;
337-
338328
let rg_configmap = build::config_map::build_rolegroup_config_map(
339329
kafka,
340330
&validated_cluster,
341331
&rolegroup_ref,
342332
validated_rg,
343333
&kafka_listeners,
344-
&pod_descriptors,
345334
)
346335
.context(BuildConfigMapSnafu)?;
347336

rust/operator-binary/src/controller/build/config_map.rs

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -12,8 +12,8 @@ use stackable_operator::{
1212
use crate::{
1313
controller::{KAFKA_CONTROLLER_NAME, ValidatedKafkaCluster, ValidatedRoleGroupConfig},
1414
crd::{
15-
JVM_SECURITY_PROPERTIES_FILE, KafkaPodDescriptor, MetadataManager,
16-
STACKABLE_LISTENER_BOOTSTRAP_DIR, STACKABLE_LISTENER_BROKER_DIR,
15+
JVM_SECURITY_PROPERTIES_FILE, MetadataManager, STACKABLE_LISTENER_BOOTSTRAP_DIR,
16+
STACKABLE_LISTENER_BROKER_DIR,
1717
listener::{KafkaListenerConfig, node_address_cmd},
1818
role::AnyConfig,
1919
v1alpha1,
@@ -66,6 +66,9 @@ pub enum Error {
6666

6767
#[snafu(display("failed to build jaas configuration file for {rolegroup}"))]
6868
BuildJaasConfig { rolegroup: String },
69+
70+
#[snafu(display("failed to build pod descriptors"))]
71+
BuildPodDescriptors { source: crate::crd::Error },
6972
}
7073

7174
/// The rolegroup [`ConfigMap`] configures the rolegroup based on the configuration given by the administrator
@@ -75,7 +78,6 @@ pub fn build_rolegroup_config_map(
7578
rolegroup: &RoleGroupRef<v1alpha1::KafkaCluster>,
7679
validated_rg: &ValidatedRoleGroupConfig,
7780
listener_config: &KafkaListenerConfig,
78-
pod_descriptors: &[KafkaPodDescriptor],
7981
) -> Result<ConfigMap, Error> {
8082
let kafka_security = &validated_cluster.kafka_security;
8183
let resolved_product_image = &validated_cluster.image;
@@ -87,6 +89,14 @@ pub fn build_rolegroup_config_map(
8789
.as_ref()
8890
.map(|auth_config| auth_config.opa_connect.clone());
8991

92+
let pod_descriptors = &kafka
93+
.pod_descriptors(
94+
None,
95+
&validated_cluster.kubernetes_cluster_info,
96+
validated_cluster.kafka_security.client_port(),
97+
)
98+
.context(BuildPodDescriptorsSnafu)?;
99+
90100
let metadata_manager = kafka
91101
.effective_metadata_manager()
92102
.context(InvalidMetadataManagerSnafu)?;

rust/operator-binary/src/controller/dereference.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@
99
//! and stays here as-is.
1010
1111
use snafu::{ResultExt, Snafu};
12-
use stackable_operator::client::Client;
12+
use stackable_operator::{client::Client, utils::cluster_info::KubernetesClusterInfo};
1313

1414
use crate::crd::{
1515
authentication::{self, ResolvedAuthenticationClasses},
@@ -33,6 +33,7 @@ type Result<T, E = Error> = std::result::Result<T, E>;
3333
pub struct DereferencedObjects {
3434
pub authentication_classes: ResolvedAuthenticationClasses,
3535
pub authorization_config: Option<KafkaAuthorizationConfig>,
36+
pub kubernetes_cluster_info: KubernetesClusterInfo,
3637
}
3738

3839
/// Fetches all Kubernetes objects referenced from the [`v1alpha1::KafkaCluster`] spec.
@@ -59,5 +60,6 @@ pub async fn dereference(
5960
Ok(DereferencedObjects {
6061
authentication_classes,
6162
authorization_config,
63+
kubernetes_cluster_info: client.kubernetes_cluster_info.clone(),
6264
})
6365
}

rust/operator-binary/src/controller/validate.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -146,6 +146,7 @@ pub fn validate(
146146
kafka_security,
147147
authorization_config: dereferenced_objects.authorization_config,
148148
role_groups,
149+
kubernetes_cluster_info: dereferenced_objects.kubernetes_cluster_info,
149150
})
150151
}
151152

0 commit comments

Comments
 (0)