Skip to content

Commit cea6ee4

Browse files
committed
fix: extract hardcoded constants
1 parent 11ca379 commit cea6ee4

4 files changed

Lines changed: 67 additions & 37 deletions

File tree

rust/operator-binary/src/controller/build/resource/statefulset.rs

Lines changed: 11 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -81,6 +81,11 @@ stackable_operator::constant!(VECTOR_CONTAINER_NAME: ContainerName = "vector");
8181
stackable_operator::constant!(VECTOR_CONFIG_VOLUME_NAME: VolumeName = "config");
8282
stackable_operator::constant!(VECTOR_LOG_VOLUME_NAME: VolumeName = "log");
8383

84+
/// Name of both the env var and the ZooKeeper discovery ConfigMap key holding the
85+
/// ZooKeeper connection string.
86+
const ZOOKEEPER_ENV_VAR_NAME: &str = "ZOOKEEPER";
87+
const POD_MANAGEMENT_POLICY_PARALLEL: &str = "Parallel";
88+
8489
#[derive(Snafu, Debug)]
8590
pub enum Error {
8691
#[snafu(display("failed to add kerberos config"))]
@@ -212,11 +217,11 @@ pub fn build_broker_rolegroup_statefulset(
212217
&validated_cluster.cluster_config.zookeeper_config_map_name
213218
{
214219
env.push(EnvVar {
215-
name: "ZOOKEEPER".to_string(),
220+
name: ZOOKEEPER_ENV_VAR_NAME.to_string(),
216221
value_from: Some(EnvVarSource {
217222
config_map_key_ref: Some(ConfigMapKeySelector {
218223
name: zookeeper_config_map_name.to_string(),
219-
key: "ZOOKEEPER".to_string(),
224+
key: ZOOKEEPER_ENV_VAR_NAME.to_string(),
220225
..ConfigMapKeySelector::default()
221226
}),
222227
..EnvVarSource::default()
@@ -410,7 +415,7 @@ pub fn build_broker_rolegroup_statefulset(
410415
.with_label(RESTART_CONTROLLER_ENABLED_LABEL.to_owned())
411416
.build(),
412417
spec: Some(StatefulSetSpec {
413-
pod_management_policy: Some("Parallel".to_string()),
418+
pod_management_policy: Some(POD_MANAGEMENT_POLICY_PARALLEL.to_string()),
414419
replicas: validated_rg.replicas.map(i32::from),
415420
selector: LabelSelector {
416421
match_labels: Some(
@@ -500,11 +505,11 @@ pub fn build_controller_rolegroup_statefulset(
500505
&validated_cluster.cluster_config.zookeeper_config_map_name
501506
{
502507
env.push(EnvVar {
503-
name: "ZOOKEEPER".to_string(),
508+
name: ZOOKEEPER_ENV_VAR_NAME.to_string(),
504509
value_from: Some(EnvVarSource {
505510
config_map_key_ref: Some(ConfigMapKeySelector {
506511
name: zookeeper_config_map_name.to_string(),
507-
key: "ZOOKEEPER".to_string(),
512+
key: ZOOKEEPER_ENV_VAR_NAME.to_string(),
508513
..ConfigMapKeySelector::default()
509514
}),
510515
..EnvVarSource::default()
@@ -635,7 +640,7 @@ pub fn build_controller_rolegroup_statefulset(
635640
.with_label(RESTART_CONTROLLER_ENABLED_LABEL.to_owned())
636641
.build(),
637642
spec: Some(StatefulSetSpec {
638-
pod_management_policy: Some("Parallel".to_string()),
643+
pod_management_policy: Some(POD_MANAGEMENT_POLICY_PARALLEL.to_string()),
639644
update_strategy: Some(StatefulSetUpdateStrategy {
640645
type_: Some("RollingUpdate".to_string()),
641646
..StatefulSetUpdateStrategy::default()

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

Lines changed: 42 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,13 @@ const OPA_TLS_MOUNT_PATH: &str = "/stackable/tls-opa";
4242
const OPA_TLS_VOLUME_NAME: &str = "tls-opa";
4343
const SSL_STORE_PASSWORD: &str = "";
4444
const SSL_STORE_TYPE_PKCS12: &str = "PKCS12";
45+
const SSL_CLIENT_AUTH_REQUIRED: &str = "required";
46+
const SASL_MECHANISM_GSSAPI: &str = "GSSAPI";
47+
const STACKABLE_KCAT_BINARY: &str = "/stackable/kcat";
48+
const PROPERTY_SECURITY_PROTOCOL: &str = "security.protocol";
49+
const PROPERTY_SASL_ENABLED_MECHANISMS: &str = "sasl.enabled.mechanisms";
50+
const PROPERTY_SASL_KERBEROS_SERVICE_NAME: &str = "sasl.kerberos.service.name";
51+
const PROPERTY_SASL_INTER_BROKER_MECHANISM: &str = "sasl.mechanism.inter.broker.protocol";
4552
const STACKABLE_TLS_KAFKA_INTERNAL_DIR: &str = "/stackable/tls-kafka-internal";
4653
const STACKABLE_TLS_KAFKA_INTERNAL_VOLUME_NAME: &str = "tls-kafka-internal";
4754
const STACKABLE_TLS_KAFKA_SERVER_DIR: &str = "/stackable/tls-kafka-server";
@@ -91,7 +98,7 @@ pub fn kcat_prober_container_commands(security: &ValidatedKafkaSecurity) -> Vec<
9198
let port = security.client_port();
9299

93100
if security.tls_client_authentication_class().is_some() {
94-
args.push("/stackable/kcat".to_string());
101+
args.push(STACKABLE_KCAT_BINARY.to_string());
95102
args.push("-b".to_string());
96103
args.push(format!("localhost:{}", port));
97104
args.extend(kcat_client_auth_ssl(STACKABLE_TLS_KCAT_DIR));
@@ -130,21 +137,21 @@ pub fn kcat_prober_container_commands(security: &ValidatedKafkaSecurity) -> Vec<
130137
)
131138
.to_string(),
132139
);
133-
bash_args.push("/stackable/kcat".to_string());
140+
bash_args.push(STACKABLE_KCAT_BINARY.to_string());
134141
bash_args.push("-b".to_string());
135142
bash_args.push("$POD_BROKER_LISTENER_ADDRESS:$POD_BROKER_LISTENER_PORT".to_string());
136143
bash_args.extend(kcat_client_sasl_ssl(STACKABLE_TLS_KCAT_DIR, service_name));
137144
bash_args.push("-L".to_string());
138145

139146
args.push(bash_args.join(" "));
140147
} else if security.tls_server_secret_class().is_some() {
141-
args.push("/stackable/kcat".to_string());
148+
args.push(STACKABLE_KCAT_BINARY.to_string());
142149
args.push("-b".to_string());
143150
args.push(format!("localhost:{}", port));
144151
args.extend(kcat_client_ssl(STACKABLE_TLS_KCAT_DIR));
145152
args.push("-L".to_string());
146153
} else {
147-
args.push("/stackable/kcat".to_string());
154+
args.push(STACKABLE_KCAT_BINARY.to_string());
148155
args.push("-b".to_string());
149156
args.push(format!("localhost:{}", port));
150157
args.push("-L".to_string());
@@ -160,10 +167,13 @@ pub fn client_properties(security: &ValidatedKafkaSecurity) -> Vec<(String, Opti
160167

161168
if security.tls_client_authentication_class().is_some() {
162169
props.push((
163-
"security.protocol".to_string(),
170+
PROPERTY_SECURITY_PROTOCOL.to_string(),
164171
Some(KafkaListenerProtocol::Ssl.to_string()),
165172
));
166-
props.push(("ssl.client.auth".to_string(), Some("required".to_string())));
173+
props.push((
174+
"ssl.client.auth".to_string(),
175+
Some(SSL_CLIENT_AUTH_REQUIRED.to_string()),
176+
));
167177
push_client_ssl_stores(&mut props, STACKABLE_TLS_KAFKA_SERVER_DIR);
168178
} else if security.has_kerberos_enabled() {
169179
// TODO: to make this configuration file usable out of the box the operator needs to be
@@ -172,21 +182,21 @@ pub fn client_properties(security: &ValidatedKafkaSecurity) -> Vec<(String, Opti
172182
// This will simplify the code and the command lines lot.
173183
// It will also make the jaas files reusable by the Kafka shell scripts.
174184
props.push((
175-
"security.protocol".to_string(),
185+
PROPERTY_SECURITY_PROTOCOL.to_string(),
176186
Some(KafkaListenerProtocol::SaslSsl.to_string()),
177187
));
178188
push_client_ssl_stores(&mut props, STACKABLE_TLS_KAFKA_SERVER_DIR);
179189
props.push((
180-
"sasl.enabled.mechanisms".to_string(),
181-
Some("GSSAPI".to_string()),
190+
PROPERTY_SASL_ENABLED_MECHANISMS.to_string(),
191+
Some(SASL_MECHANISM_GSSAPI.to_string()),
182192
));
183193
props.push((
184-
"sasl.kerberos.service.name".to_string(),
194+
PROPERTY_SASL_KERBEROS_SERVICE_NAME.to_string(),
185195
Some(KafkaRole::Broker.kerberos_service_name().to_string()),
186196
));
187197
props.push((
188-
"sasl.mechanism.inter.broker.protocol".to_string(),
189-
Some("GSSAPI".to_string()),
198+
PROPERTY_SASL_INTER_BROKER_MECHANISM.to_string(),
199+
Some(SASL_MECHANISM_GSSAPI.to_string()),
190200
));
191201
props.push((
192202
"sasl.jaas.config".to_string(),
@@ -197,13 +207,13 @@ pub fn client_properties(security: &ValidatedKafkaSecurity) -> Vec<(String, Opti
197207
realm="$KERBEROS_REALM"))));
198208
} else if security.tls_server_secret_class().is_some() {
199209
props.push((
200-
"security.protocol".to_string(),
210+
PROPERTY_SECURITY_PROTOCOL.to_string(),
201211
Some(KafkaListenerProtocol::Ssl.to_string()),
202212
));
203213
push_client_ssl_truststore(&mut props, STACKABLE_TLS_KAFKA_SERVER_DIR);
204214
} else {
205215
props.push((
206-
"security.protocol".to_string(),
216+
PROPERTY_SECURITY_PROTOCOL.to_string(),
207217
Some(KafkaListenerProtocol::Plaintext.to_string()),
208218
));
209219
}
@@ -417,7 +427,7 @@ pub fn broker_config_settings(security: &ValidatedKafkaSecurity) -> BTreeMap<Str
417427
// client auth required
418428
config.insert(
419429
KafkaListenerName::Client.listener_ssl_client_auth(),
420-
"required".to_string(),
430+
SSL_CLIENT_AUTH_REQUIRED.to_string(),
421431
);
422432
}
423433
}
@@ -429,14 +439,17 @@ pub fn broker_config_settings(security: &ValidatedKafkaSecurity) -> BTreeMap<Str
429439
&KafkaListenerName::Bootstrap,
430440
STACKABLE_TLS_KAFKA_SERVER_DIR,
431441
);
432-
config.insert("sasl.enabled.mechanisms".to_string(), "GSSAPI".to_string());
433442
config.insert(
434-
"sasl.kerberos.service.name".to_string(),
443+
PROPERTY_SASL_ENABLED_MECHANISMS.to_string(),
444+
SASL_MECHANISM_GSSAPI.to_string(),
445+
);
446+
config.insert(
447+
PROPERTY_SASL_KERBEROS_SERVICE_NAME.to_string(),
435448
KafkaRole::Broker.kerberos_service_name().to_string(),
436449
);
437450
config.insert(
438-
"sasl.mechanism.inter.broker.protocol".to_string(),
439-
"GSSAPI".to_string(),
451+
PROPERTY_SASL_INTER_BROKER_MECHANISM.to_string(),
452+
SASL_MECHANISM_GSSAPI.to_string(),
440453
);
441454
tracing::debug!("Kerberos configs added: [{:#?}]", config);
442455
}
@@ -458,7 +471,7 @@ pub fn broker_config_settings(security: &ValidatedKafkaSecurity) -> BTreeMap<Str
458471
// client auth required
459472
config.insert(
460473
KafkaListenerName::Internal.listener_ssl_client_auth(),
461-
"required".to_string(),
474+
SSL_CLIENT_AUTH_REQUIRED.to_string(),
462475
);
463476
}
464477

@@ -518,21 +531,24 @@ pub fn controller_config_settings(security: &ValidatedKafkaSecurity) -> BTreeMap
518531
// client auth required
519532
config.insert(
520533
KafkaListenerName::Controller.listener_ssl_client_auth(),
521-
"required".to_string(),
534+
SSL_CLIENT_AUTH_REQUIRED.to_string(),
522535
);
523536
}
524537
}
525538

526539
// Kerberos
527540
if security.has_kerberos_enabled() {
528-
config.insert("sasl.enabled.mechanisms".to_string(), "GSSAPI".to_string());
529541
config.insert(
530-
"sasl.kerberos.service.name".to_string(),
542+
PROPERTY_SASL_ENABLED_MECHANISMS.to_string(),
543+
SASL_MECHANISM_GSSAPI.to_string(),
544+
);
545+
config.insert(
546+
PROPERTY_SASL_KERBEROS_SERVICE_NAME.to_string(),
531547
KafkaRole::Controller.kerberos_service_name().to_string(),
532548
);
533549
config.insert(
534-
"sasl.mechanism.inter.broker.protocol".to_string(),
535-
"GSSAPI".to_string(),
550+
PROPERTY_SASL_INTER_BROKER_MECHANISM.to_string(),
551+
SASL_MECHANISM_GSSAPI.to_string(),
536552
);
537553
tracing::debug!("Kerberos configs added: [{:#?}]", config);
538554
}
@@ -631,7 +647,7 @@ fn kcat_client_sasl_ssl(cert_directory: &str, service_name: &str) -> Vec<String>
631647
"-X".to_string(),
632648
"sasl.kerberos.keytab=/stackable/kerberos/keytab".to_string(),
633649
"-X".to_string(),
634-
"sasl.mechanism=GSSAPI".to_string(),
650+
format!("sasl.mechanism={SASL_MECHANISM_GSSAPI}"),
635651
"-X".to_string(),
636652
format!("sasl.kerberos.service.name={service_name}"),
637653
"-X".to_string(),

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

Lines changed: 13 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,13 @@ use strum::EnumString;
77

88
pub(crate) const LISTENER_LOCAL_ADDRESS: &str = "0.0.0.0";
99

10+
// Layout of the listener volume mounted by the listener operator: the default address is
11+
// exposed under `<dir>/default-address/address` and per-port values under
12+
// `<dir>/default-address/ports/<port-name>`.
13+
const LISTENER_DEFAULT_ADDRESS_DIR: &str = "default-address";
14+
const LISTENER_ADDRESS_FILE: &str = "address";
15+
const LISTENER_PORTS_DIR: &str = "ports";
16+
1017
#[derive(strum::Display, Debug, EnumString)]
1118
pub enum KafkaListenerProtocol {
1219
/// Unencrypted and unauthenticated HTTP connections
@@ -178,17 +185,19 @@ impl Display for KafkaListener {
178185
}
179186

180187
pub fn node_address_cmd_env(directory: &str) -> String {
181-
format!("$(cat {directory}/default-address/address)")
188+
format!("$(cat {directory}/{LISTENER_DEFAULT_ADDRESS_DIR}/{LISTENER_ADDRESS_FILE})")
182189
}
183190

184191
pub fn node_port_cmd_env(directory: &str, port_name: &str) -> String {
185-
format!("$(cat {directory}/default-address/ports/{port_name})")
192+
format!("$(cat {directory}/{LISTENER_DEFAULT_ADDRESS_DIR}/{LISTENER_PORTS_DIR}/{port_name})")
186193
}
187194

188195
pub fn node_address_cmd(directory: &str) -> String {
189-
format!("${{file:UTF-8:{directory}/default-address/address}}")
196+
format!("${{file:UTF-8:{directory}/{LISTENER_DEFAULT_ADDRESS_DIR}/{LISTENER_ADDRESS_FILE}}}")
190197
}
191198

192199
pub fn node_port_cmd(directory: &str, port_name: &str) -> String {
193-
format!("${{file:UTF-8:{directory}/default-address/ports/{port_name}}}")
200+
format!(
201+
"${{file:UTF-8:{directory}/{LISTENER_DEFAULT_ADDRESS_DIR}/{LISTENER_PORTS_DIR}/{port_name}}}"
202+
)
194203
}

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -280,7 +280,7 @@ impl HasStatusCondition for v1alpha1::KafkaCluster {
280280
impl v1alpha1::KafkaCluster {
281281
pub fn effective_metadata_manager(&self) -> Result<MetadataManager, Error> {
282282
match &self.spec.cluster_config.metadata_manager {
283-
Some(manager) => match manager.clone() {
283+
Some(manager) => match manager {
284284
MetadataManager::ZooKeeper => {
285285
if !self.spec.image.product_version().starts_with("3.") {
286286
Err(Error::Kafka4RequiresKraftMetadataManager)

0 commit comments

Comments
 (0)