Skip to content

Commit 8ae1c5f

Browse files
feat: Use separate ConfigOverrides structs for brokers and controllers
1 parent 572179d commit 8ae1c5f

4 files changed

Lines changed: 112 additions & 110 deletions

File tree

extra/crds.yaml

Lines changed: 0 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -549,16 +549,6 @@ spec:
549549
used by `HashMap<String, String>`.
550550
nullable: true
551551
type: object
552-
controller.properties:
553-
additionalProperties:
554-
type: string
555-
description: |-
556-
Flat key-value overrides for `*.properties`, Hadoop XML, etc.
557-
558-
This is backwards-compatible with the existing flat key-value YAML format
559-
used by `HashMap<String, String>`.
560-
nullable: true
561-
type: object
562552
security.properties:
563553
additionalProperties:
564554
type: string
@@ -1176,16 +1166,6 @@ spec:
11761166
used by `HashMap<String, String>`.
11771167
nullable: true
11781168
type: object
1179-
controller.properties:
1180-
additionalProperties:
1181-
type: string
1182-
description: |-
1183-
Flat key-value overrides for `*.properties`, Hadoop XML, etc.
1184-
1185-
This is backwards-compatible with the existing flat key-value YAML format
1186-
used by `HashMap<String, String>`.
1187-
nullable: true
1188-
type: object
11891169
security.properties:
11901170
additionalProperties:
11911171
type: string
@@ -1804,16 +1784,6 @@ spec:
18041784
and consult the operator specific usage guide documentation for details on the
18051785
available config files and settings for the specific product.
18061786
properties:
1807-
broker.properties:
1808-
additionalProperties:
1809-
type: string
1810-
description: |-
1811-
Flat key-value overrides for `*.properties`, Hadoop XML, etc.
1812-
1813-
This is backwards-compatible with the existing flat key-value YAML format
1814-
used by `HashMap<String, String>`.
1815-
nullable: true
1816-
type: object
18171787
controller.properties:
18181788
additionalProperties:
18191789
type: string
@@ -2263,16 +2233,6 @@ spec:
22632233
and consult the operator specific usage guide documentation for details on the
22642234
available config files and settings for the specific product.
22652235
properties:
2266-
broker.properties:
2267-
additionalProperties:
2268-
type: string
2269-
description: |-
2270-
Flat key-value overrides for `*.properties`, Hadoop XML, etc.
2271-
2272-
This is backwards-compatible with the existing flat key-value YAML format
2273-
used by `HashMap<String, String>`.
2274-
nullable: true
2275-
type: object
22762236
controller.properties:
22772237
additionalProperties:
22782238
type: string

rust/operator-binary/src/crd/mod.rs

Lines changed: 76 additions & 44 deletions
Original file line numberDiff line numberDiff line change
@@ -60,50 +60,6 @@ pub const STACKABLE_CONFIG_DIR: &str = "/stackable/config";
6060
pub const STACKABLE_KERBEROS_DIR: &str = "/stackable/kerberos";
6161
pub const STACKABLE_KERBEROS_KRB5_PATH: &str = "/stackable/kerberos/krb5.conf";
6262

63-
pub type BrokerRole =
64-
Role<BrokerConfigFragment, KafkaConfigOverrides, GenericRoleConfig, JavaCommonConfig>;
65-
pub type ControllerRole =
66-
Role<ControllerConfigFragment, KafkaConfigOverrides, GenericRoleConfig, JavaCommonConfig>;
67-
68-
#[derive(Clone, Debug, Default, Deserialize, JsonSchema, PartialEq, Serialize)]
69-
#[serde(rename_all = "camelCase")]
70-
pub struct KafkaConfigOverrides {
71-
#[serde(
72-
default,
73-
rename = "broker.properties",
74-
skip_serializing_if = "Option::is_none"
75-
)]
76-
pub broker_properties: Option<KeyValueConfigOverrides>,
77-
78-
#[serde(
79-
default,
80-
rename = "controller.properties",
81-
skip_serializing_if = "Option::is_none"
82-
)]
83-
pub controller_properties: Option<KeyValueConfigOverrides>,
84-
85-
#[serde(
86-
default,
87-
rename = "security.properties",
88-
skip_serializing_if = "Option::is_none"
89-
)]
90-
pub security_properties: Option<KeyValueConfigOverrides>,
91-
}
92-
93-
impl KeyValueOverridesProvider for KafkaConfigOverrides {
94-
fn get_key_value_overrides(&self, file: &str) -> BTreeMap<String, Option<String>> {
95-
let field = match file {
96-
role::broker::BROKER_PROPERTIES_FILE => self.broker_properties.as_ref(),
97-
role::controller::CONTROLLER_PROPERTIES_FILE => self.controller_properties.as_ref(),
98-
JVM_SECURITY_PROPERTIES_FILE => self.security_properties.as_ref(),
99-
_ => None,
100-
};
101-
field
102-
.map(KeyValueConfigOverrides::as_product_config_overrides)
103-
.unwrap_or_default()
104-
}
105-
}
106-
10763
#[derive(Snafu, Debug)]
10864
pub enum Error {
10965
#[snafu(display(
@@ -138,6 +94,20 @@ pub enum Error {
13894
},
13995
}
14096

97+
pub type BrokerRole = Role<
98+
BrokerConfigFragment,
99+
v1alpha1::KafkaBrokerConfigOverrides,
100+
GenericRoleConfig,
101+
JavaCommonConfig,
102+
>;
103+
104+
pub type ControllerRole = Role<
105+
ControllerConfigFragment,
106+
v1alpha1::KafkaControllerConfigOverrides,
107+
GenericRoleConfig,
108+
JavaCommonConfig,
109+
>;
110+
141111
#[versioned(
142112
version(name = "v1alpha1"),
143113
crates(
@@ -269,6 +239,42 @@ pub mod versioned {
269239
#[serde(skip_serializing_if = "Option::is_none")]
270240
pub broker_id_pod_config_map_name: Option<String>,
271241
}
242+
243+
#[derive(Clone, Debug, Default, Deserialize, JsonSchema, PartialEq, Serialize)]
244+
#[serde(rename_all = "camelCase")]
245+
pub struct KafkaBrokerConfigOverrides {
246+
#[serde(
247+
default,
248+
rename = "broker.properties",
249+
skip_serializing_if = "Option::is_none"
250+
)]
251+
pub broker_properties: Option<KeyValueConfigOverrides>,
252+
253+
#[serde(
254+
default,
255+
rename = "security.properties",
256+
skip_serializing_if = "Option::is_none"
257+
)]
258+
pub security_properties: Option<KeyValueConfigOverrides>,
259+
}
260+
261+
#[derive(Clone, Debug, Default, Deserialize, JsonSchema, PartialEq, Serialize)]
262+
#[serde(rename_all = "camelCase")]
263+
pub struct KafkaControllerConfigOverrides {
264+
#[serde(
265+
default,
266+
rename = "controller.properties",
267+
skip_serializing_if = "Option::is_none"
268+
)]
269+
pub controller_properties: Option<KeyValueConfigOverrides>,
270+
271+
#[serde(
272+
default,
273+
rename = "security.properties",
274+
skip_serializing_if = "Option::is_none"
275+
)]
276+
pub security_properties: Option<KeyValueConfigOverrides>,
277+
}
272278
}
273279

274280
impl Default for v1alpha1::KafkaClusterConfig {
@@ -285,6 +291,32 @@ impl Default for v1alpha1::KafkaClusterConfig {
285291
}
286292
}
287293

294+
impl KeyValueOverridesProvider for v1alpha1::KafkaBrokerConfigOverrides {
295+
fn get_key_value_overrides(&self, file: &str) -> BTreeMap<String, Option<String>> {
296+
let field = match file {
297+
role::broker::BROKER_PROPERTIES_FILE => self.broker_properties.as_ref(),
298+
JVM_SECURITY_PROPERTIES_FILE => self.security_properties.as_ref(),
299+
_ => None,
300+
};
301+
field
302+
.map(KeyValueConfigOverrides::as_product_config_overrides)
303+
.unwrap_or_default()
304+
}
305+
}
306+
307+
impl KeyValueOverridesProvider for v1alpha1::KafkaControllerConfigOverrides {
308+
fn get_key_value_overrides(&self, file: &str) -> BTreeMap<String, Option<String>> {
309+
let field = match file {
310+
role::controller::CONTROLLER_PROPERTIES_FILE => self.controller_properties.as_ref(),
311+
JVM_SECURITY_PROPERTIES_FILE => self.security_properties.as_ref(),
312+
_ => None,
313+
};
314+
field
315+
.map(KeyValueConfigOverrides::as_product_config_overrides)
316+
.unwrap_or_default()
317+
}
318+
}
319+
288320
impl HasStatusCondition for v1alpha1::KafkaCluster {
289321
fn conditions(&self) -> Vec<ClusterCondition> {
290322
match &self.status {

rust/operator-binary/src/crd/role/mod.rs

Lines changed: 19 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -22,13 +22,10 @@ use strum::{Display, EnumIter, EnumString, IntoEnumIterator};
2222

2323
use crate::{
2424
config::jvm::{construct_heap_jvm_args, construct_non_heap_jvm_args},
25-
crd::{
26-
KafkaConfigOverrides,
27-
role::{
28-
broker::{BROKER_PROPERTIES_FILE, BrokerConfig, BrokerConfigFragment},
29-
commons::{CommonConfig, Storage},
30-
controller::{CONTROLLER_PROPERTIES_FILE, ControllerConfig},
31-
},
25+
crd::role::{
26+
broker::{BROKER_PROPERTIES_FILE, BrokerConfig, BrokerConfigFragment},
27+
commons::{CommonConfig, Storage},
28+
controller::{CONTROLLER_PROPERTIES_FILE, ControllerConfig},
3229
},
3330
v1alpha1,
3431
};
@@ -229,16 +226,17 @@ impl KafkaRole {
229226
rolegroup: &str,
230227
) -> Result<String, Error> {
231228
match self {
232-
Self::Broker => {
233-
construct_non_heap_jvm_args::<BrokerConfigFragment, KafkaConfigOverrides>(
234-
merged_config,
235-
kafka.broker_role().with_context(|_| MissingRoleSnafu {
236-
role: self.to_string(),
237-
})?,
238-
rolegroup,
239-
)
240-
.context(ConstructJvmArgumentsSnafu)
241-
}
229+
Self::Broker => construct_non_heap_jvm_args::<
230+
BrokerConfigFragment,
231+
v1alpha1::KafkaBrokerConfigOverrides,
232+
>(
233+
merged_config,
234+
kafka.broker_role().with_context(|_| MissingRoleSnafu {
235+
role: self.to_string(),
236+
})?,
237+
rolegroup,
238+
)
239+
.context(ConstructJvmArgumentsSnafu),
242240
Self::Controller => construct_non_heap_jvm_args(
243241
merged_config,
244242
kafka.controller_role().with_context(|_| MissingRoleSnafu {
@@ -257,7 +255,10 @@ impl KafkaRole {
257255
rolegroup: &str,
258256
) -> Result<String, Error> {
259257
match self {
260-
Self::Broker => construct_heap_jvm_args::<BrokerConfigFragment, KafkaConfigOverrides>(
258+
Self::Broker => construct_heap_jvm_args::<
259+
BrokerConfigFragment,
260+
v1alpha1::KafkaBrokerConfigOverrides,
261+
>(
261262
merged_config,
262263
kafka.broker_role().with_context(|_| MissingRoleSnafu {
263264
role: self.to_string(),

rust/operator-binary/src/kafka_controller.rs

Lines changed: 17 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -546,9 +546,9 @@ fn validated_product_config(
546546
product_version: &str,
547547
product_config: &ProductConfigManager,
548548
) -> Result<ValidatedRoleConfigByPropertyKind, Error> {
549-
let mut roles = HashMap::new();
549+
let mut role_config = HashMap::new();
550550

551-
roles.insert(
551+
let broker_role = [(
552552
KafkaRole::Broker.to_string(),
553553
(
554554
vec![
@@ -564,11 +564,17 @@ fn validated_product_config(
564564
})?
565565
.erase(),
566566
),
567-
);
567+
)]
568+
.into();
569+
570+
let broker_role_config =
571+
transform_all_roles_to_config(kafka, &broker_role).context(GenerateProductConfigSnafu)?;
572+
573+
role_config.extend(broker_role_config);
568574

569575
// TODO: need this if because controller_role() raises an error
570576
if kafka.spec.controllers.is_some() {
571-
roles.insert(
577+
let controller_role = [(
572578
KafkaRole::Controller.to_string(),
573579
(
574580
vec![
@@ -584,11 +590,14 @@ fn validated_product_config(
584590
})?
585591
.erase(),
586592
),
587-
);
588-
}
593+
)]
594+
.into();
595+
596+
let controller_role_config = transform_all_roles_to_config(kafka, &controller_role)
597+
.context(GenerateProductConfigSnafu)?;
589598

590-
let role_config =
591-
transform_all_roles_to_config(kafka, &roles).context(GenerateProductConfigSnafu)?;
599+
role_config.extend(controller_role_config);
600+
}
592601

593602
validate_all_roles_and_groups_config(
594603
product_version,

0 commit comments

Comments
 (0)