Skip to content

Commit af0ceba

Browse files
committed
move resolution of metadata_manager and pod descriptors to validate stage
1 parent 4eb786b commit af0ceba

6 files changed

Lines changed: 45 additions & 58 deletions

File tree

rust/operator-binary/src/controller.rs

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

@@ -32,7 +31,7 @@ mod validate;
3231

3332
use crate::{
3433
crd::{
35-
self, APP_NAME, KafkaClusterStatus, OPERATOR_NAME,
34+
self, APP_NAME, KafkaClusterStatus, KafkaPodDescriptor, MetadataManager, OPERATOR_NAME,
3635
authorization::KafkaAuthorizationConfig,
3736
listener::get_kafka_listener_config,
3837
role::{AnyConfig, KafkaRole},
@@ -219,7 +218,8 @@ pub struct ValidatedKafkaCluster {
219218
// classes — rejected because nothing downstream needs them beyond kafka_security.
220219
pub authorization_config: Option<KafkaAuthorizationConfig>,
221220
pub role_groups: BTreeMap<KafkaRole, BTreeMap<String, ValidatedRoleGroupConfig>>,
222-
pub kubernetes_cluster_info: KubernetesClusterInfo,
221+
pub pod_descriptors: Vec<KafkaPodDescriptor>,
222+
pub metadata_manager: MetadataManager,
223223
}
224224

225225
pub struct ValidatedRoleGroupConfig {

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

Lines changed: 12 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -58,17 +58,14 @@ pub enum Error {
5858
rolegroup: RoleGroupRef<v1alpha1::KafkaCluster>,
5959
},
6060

61-
#[snafu(display("failed to build properties for {rolegroup}"))]
62-
BuildProperties {
63-
source: crate::controller::build::properties::Error,
64-
rolegroup: RoleGroupRef<v1alpha1::KafkaCluster>,
65-
},
66-
6761
#[snafu(display("failed to build jaas configuration file for {rolegroup}"))]
6862
BuildJaasConfig { rolegroup: String },
6963

7064
#[snafu(display("failed to build pod descriptors"))]
7165
BuildPodDescriptors { source: crate::crd::Error },
66+
67+
#[snafu(display("no Kraft controllers found to build"))]
68+
NoKraftControllersFound,
7269
}
7370

7471
/// The rolegroup [`ConfigMap`] configures the rolegroup based on the configuration given by the administrator
@@ -89,25 +86,19 @@ pub fn build_rolegroup_config_map(
8986
.as_ref()
9087
.map(|auth_config| auth_config.opa_connect.clone());
9188

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)?;
89+
let kraft_mode = validated_cluster.metadata_manager == MetadataManager::KRaft;
9990

100-
let metadata_manager = kafka
101-
.effective_metadata_manager()
102-
.context(InvalidMetadataManagerSnafu)?;
91+
if kraft_mode && validated_cluster.pod_descriptors.is_empty() {
92+
return NoKraftControllersFoundSnafu.fail();
93+
}
10394

10495
let kafka_config = match &validated_rg.merged_config {
10596
AnyConfig::Broker(_) => crate::controller::build::properties::broker_properties::build(
10697
kafka_security,
10798
listener_config,
108-
pod_descriptors,
99+
&validated_cluster.pod_descriptors,
109100
opa_connect.as_deref(),
110-
metadata_manager == MetadataManager::KRaft,
101+
kraft_mode,
111102
kafka
112103
.spec
113104
.cluster_config
@@ -119,15 +110,12 @@ pub fn build_rolegroup_config_map(
119110
crate::controller::build::properties::controller_properties::build(
120111
kafka_security,
121112
listener_config,
122-
pod_descriptors,
123-
metadata_manager == MetadataManager::KRaft,
113+
&validated_cluster.pod_descriptors,
114+
kraft_mode,
124115
config_overrides,
125116
)
126117
}
127-
}
128-
.with_context(|_| BuildPropertiesSnafu {
129-
rolegroup: rolegroup.clone(),
130-
})?;
118+
};
131119

132120
let kafka_config = kafka_config
133121
.into_iter()

rust/operator-binary/src/controller/build/properties/broker_properties.rs

Lines changed: 4 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,6 @@
11
use std::collections::BTreeMap;
22

3-
use snafu::OptionExt;
4-
5-
use super::{Error, NoKraftControllersFoundSnafu, kraft_controllers};
3+
use super::kraft_controllers;
64
use crate::{
75
crd::{
86
KafkaPodDescriptor,
@@ -25,7 +23,7 @@ pub fn build(
2523
kraft_mode: bool,
2624
disable_broker_id_generation: bool,
2725
overrides: BTreeMap<String, String>,
28-
) -> Result<BTreeMap<String, String>, Error> {
26+
) -> BTreeMap<String, String> {
2927
let kraft_controllers = kraft_controllers(pod_descriptors);
3028

3129
let mut result = BTreeMap::from([
@@ -49,7 +47,7 @@ pub fn build(
4947
]);
5048

5149
if kraft_mode {
52-
let kraft_controllers = kraft_controllers.context(NoKraftControllersFoundSnafu)?;
50+
let kraft_controllers = kraft_controllers.join(",");
5351

5452
// Running in KRaft mode
5553
result.extend([
@@ -114,5 +112,5 @@ pub fn build(
114112
result.extend(graceful_shutdown_config_properties());
115113
result.extend(overrides);
116114

117-
Ok(result)
115+
result
118116
}

rust/operator-binary/src/controller/build/properties/controller_properties.rs

Lines changed: 4 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,6 @@
11
use std::collections::BTreeMap;
22

3-
use snafu::OptionExt;
4-
5-
use super::{Error, NoKraftControllersFoundSnafu, kraft_controllers};
3+
use super::kraft_controllers;
64
use crate::{
75
crd::{
86
KafkaPodDescriptor,
@@ -22,9 +20,8 @@ pub fn build(
2220
pod_descriptors: &[KafkaPodDescriptor],
2321
kraft_mode: bool,
2422
overrides: BTreeMap<String, String>,
25-
) -> Result<BTreeMap<String, String>, Error> {
26-
let kraft_controllers =
27-
kraft_controllers(pod_descriptors).context(NoKraftControllersFoundSnafu)?;
23+
) -> BTreeMap<String, String> {
24+
let kraft_controllers = kraft_controllers(pod_descriptors).join(",");
2825

2926
let mut result = BTreeMap::from([
3027
(
@@ -72,5 +69,5 @@ pub fn build(
7269
result.extend(graceful_shutdown_config_properties());
7370
result.extend(overrides);
7471

75-
Ok(result)
72+
result
7673
}

rust/operator-binary/src/controller/build/properties/mod.rs

Lines changed: 2 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -3,18 +3,10 @@
33
pub mod broker_properties;
44
pub mod controller_properties;
55

6-
use snafu::Snafu;
7-
86
use crate::crd::{KafkaPodDescriptor, role::KafkaRole};
97

10-
#[derive(Snafu, Debug)]
11-
pub enum Error {
12-
#[snafu(display("no Kraft controllers found to build"))]
13-
NoKraftControllersFound,
14-
}
15-
16-
pub(crate) fn kraft_controllers(pod_descriptors: &[KafkaPodDescriptor]) -> Option<String> {
17-
let result = pod_descriptors
8+
pub(crate) fn kraft_controllers(pod_descriptors: &[KafkaPodDescriptor]) -> Vec<String> {
9+
pod_descriptors
1810
.iter()
1911
.filter(|pd| pd.role == KafkaRole::Controller.to_string())
2012
.map(|desc| {
@@ -25,11 +17,4 @@ pub(crate) fn kraft_controllers(pod_descriptors: &[KafkaPodDescriptor]) -> Optio
2517
)
2618
})
2719
.collect::<Vec<String>>()
28-
.join(",");
29-
30-
if result.is_empty() {
31-
None
32-
} else {
33-
Some(result)
34-
}
3520
}

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

Lines changed: 20 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,12 @@ pub enum Error {
3939

4040
#[snafu(display("failed to resolve merged config for rolegroup"))]
4141
ResolveMergedConfig { source: crate::crd::role::Error },
42+
43+
#[snafu(display("failed to build pod descriptors"))]
44+
BuildPodDescriptors { source: crate::crd::Error },
45+
46+
#[snafu(display("invalid metadata manager"))]
47+
InvalidMetadataManager { source: crate::crd::Error },
4248
}
4349

4450
type Result<T, E = Error> = std::result::Result<T, E>;
@@ -141,12 +147,25 @@ pub fn validate(
141147
role_groups.insert(KafkaRole::Controller, controller_groups);
142148
}
143149

150+
let pod_descriptors = kafka
151+
.pod_descriptors(
152+
None,
153+
&dereferenced_objects.kubernetes_cluster_info,
154+
kafka_security.client_port(),
155+
)
156+
.context(BuildPodDescriptorsSnafu)?;
157+
158+
let metadata_manager = kafka
159+
.effective_metadata_manager()
160+
.context(InvalidMetadataManagerSnafu)?;
161+
144162
Ok(ValidatedKafkaCluster {
145163
image,
146164
kafka_security,
147165
authorization_config: dereferenced_objects.authorization_config,
148166
role_groups,
149-
kubernetes_cluster_info: dereferenced_objects.kubernetes_cluster_info,
167+
pod_descriptors,
168+
metadata_manager,
150169
})
151170
}
152171

0 commit comments

Comments
 (0)