Skip to content

Commit 75a0519

Browse files
adwk67claude
andcommitted
refactor: extract per-file kafka .properties builders
Splits server_properties_file into controller/build/properties/{broker,controller}_properties builders (base map + security settings + graceful-shutdown + user overrides), wired into resource/configmap.rs by role. The property assembly was moved verbatim, so the rendered broker.properties/controller.properties are unchanged (18 tests pass; byte parity to be confirmed via the kuttl ConfigMap snapshot). Override input stays BTreeMap<String,String>; no product-config removed yet (later increment). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
1 parent 9225150 commit 75a0519

6 files changed

Lines changed: 273 additions & 216 deletions

File tree

rust/operator-binary/src/controller.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ use stackable_operator::{
2626
};
2727
use strum::{EnumDiscriminants, IntoStaticStr};
2828

29+
pub(crate) mod build;
2930
mod dereference;
3031
mod validate;
3132

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,3 @@
1+
//! Builders that assemble Kubernetes resources for kafka rolegroups.
2+
3+
pub mod properties;
Lines changed: 119 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,119 @@
1+
use std::collections::BTreeMap;
2+
3+
use snafu::OptionExt;
4+
5+
use crate::{
6+
crd::{
7+
KafkaPodDescriptor,
8+
listener::{KafkaListenerConfig, KafkaListenerName},
9+
role::{
10+
KAFKA_ADVERTISED_LISTENERS, KAFKA_BROKER_ID,
11+
KAFKA_CONTROLLER_QUORUM_BOOTSTRAP_SERVERS, KAFKA_LISTENER_SECURITY_PROTOCOL_MAP,
12+
KAFKA_LISTENERS, KAFKA_LOG_DIRS, KAFKA_NODE_ID, KAFKA_PROCESS_ROLES, KafkaRole,
13+
},
14+
security::KafkaTlsSecurity,
15+
},
16+
operations::graceful_shutdown::graceful_shutdown_config_properties,
17+
};
18+
19+
use super::{Error, NoKraftControllersFoundSnafu, kraft_controllers};
20+
21+
pub fn build(
22+
kafka_security: &KafkaTlsSecurity,
23+
listener_config: &KafkaListenerConfig,
24+
pod_descriptors: &[KafkaPodDescriptor],
25+
opa_connect_string: Option<&str>,
26+
kraft_mode: bool,
27+
disable_broker_id_generation: bool,
28+
overrides: BTreeMap<String, String>,
29+
) -> Result<BTreeMap<String, String>, Error> {
30+
let kraft_controllers = kraft_controllers(pod_descriptors);
31+
32+
let mut result = BTreeMap::from([
33+
(
34+
KAFKA_LOG_DIRS.to_string(),
35+
"/stackable/data/topicdata".to_string(),
36+
),
37+
(KAFKA_LISTENERS.to_string(), listener_config.listeners()),
38+
(
39+
KAFKA_ADVERTISED_LISTENERS.to_string(),
40+
listener_config.advertised_listeners(),
41+
),
42+
(
43+
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP.to_string(),
44+
listener_config.listener_security_protocol_map(),
45+
),
46+
(
47+
"inter.broker.listener.name".to_string(),
48+
KafkaListenerName::Internal.to_string(),
49+
),
50+
]);
51+
52+
if kraft_mode {
53+
let kraft_controllers = kraft_controllers.context(NoKraftControllersFoundSnafu)?;
54+
55+
// Running in KRaft mode
56+
result.extend([
57+
(
58+
"broker.id.generation.enable".to_string(),
59+
"false".to_string(),
60+
),
61+
(KAFKA_NODE_ID.to_string(), "${env:REPLICA_ID}".to_string()),
62+
(
63+
KAFKA_PROCESS_ROLES.to_string(),
64+
KafkaRole::Broker.to_string(),
65+
),
66+
(
67+
"controller.listener.names".to_string(),
68+
KafkaListenerName::Controller.to_string(),
69+
),
70+
(
71+
KAFKA_CONTROLLER_QUORUM_BOOTSTRAP_SERVERS.to_string(),
72+
kraft_controllers.clone(),
73+
),
74+
]);
75+
} else {
76+
// Running with ZooKeeper enabled
77+
result.extend([(
78+
"zookeeper.connect".to_string(),
79+
"${env:ZOOKEEPER}".to_string(),
80+
)]);
81+
// We are in zookeeper mode and the user has defined a broker id mapping
82+
// so we disable automatic id generation.
83+
// This check ensures that existing clusters running in ZooKeeper mode do not
84+
// suddenly break after the introduction of this change.
85+
if disable_broker_id_generation {
86+
result.extend([
87+
(
88+
"broker.id.generation.enable".to_string(),
89+
"false".to_string(),
90+
),
91+
(KAFKA_BROKER_ID.to_string(), "${env:REPLICA_ID}".to_string()),
92+
]);
93+
}
94+
}
95+
96+
// Enable OPA authorization
97+
if opa_connect_string.is_some() {
98+
result.extend([
99+
(
100+
"authorizer.class.name".to_string(),
101+
"org.openpolicyagent.kafka.OpaAuthorizer".to_string(),
102+
),
103+
(
104+
"opa.authorizer.metrics.enabled".to_string(),
105+
"true".to_string(),
106+
),
107+
(
108+
"opa.authorizer.url".to_string(),
109+
opa_connect_string.unwrap_or_default().to_string(),
110+
),
111+
]);
112+
}
113+
114+
result.extend(kafka_security.broker_config_settings());
115+
result.extend(graceful_shutdown_config_properties());
116+
result.extend(overrides);
117+
118+
Ok(result)
119+
}
Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,76 @@
1+
use std::collections::BTreeMap;
2+
3+
use snafu::OptionExt;
4+
5+
use crate::{
6+
crd::{
7+
KafkaPodDescriptor,
8+
listener::{KafkaListenerConfig, KafkaListenerName},
9+
role::{
10+
KAFKA_CONTROLLER_QUORUM_BOOTSTRAP_SERVERS, KAFKA_LISTENER_SECURITY_PROTOCOL_MAP,
11+
KAFKA_LISTENERS, KAFKA_LOG_DIRS, KAFKA_NODE_ID, KAFKA_PROCESS_ROLES, KafkaRole,
12+
},
13+
security::KafkaTlsSecurity,
14+
},
15+
operations::graceful_shutdown::graceful_shutdown_config_properties,
16+
};
17+
18+
use super::{Error, NoKraftControllersFoundSnafu, kraft_controllers};
19+
20+
pub fn build(
21+
kafka_security: &KafkaTlsSecurity,
22+
listener_config: &KafkaListenerConfig,
23+
pod_descriptors: &[KafkaPodDescriptor],
24+
kraft_mode: bool,
25+
overrides: BTreeMap<String, String>,
26+
) -> Result<BTreeMap<String, String>, Error> {
27+
let kraft_controllers = kraft_controllers(pod_descriptors).context(NoKraftControllersFoundSnafu)?;
28+
29+
let mut result = BTreeMap::from([
30+
(
31+
KAFKA_LOG_DIRS.to_string(),
32+
"/stackable/data/kraft".to_string(),
33+
),
34+
(KAFKA_PROCESS_ROLES.to_string(), KafkaRole::Controller.to_string()),
35+
(
36+
"controller.listener.names".to_string(),
37+
KafkaListenerName::Controller.to_string(),
38+
),
39+
(
40+
KAFKA_NODE_ID.to_string(),
41+
"${env:REPLICA_ID}".to_string(),
42+
),
43+
(
44+
KAFKA_CONTROLLER_QUORUM_BOOTSTRAP_SERVERS.to_string(),
45+
kraft_controllers.clone(),
46+
),
47+
(
48+
KAFKA_LISTENERS.to_string(),
49+
"CONTROLLER://${env:POD_NAME}.${env:ROLEGROUP_HEADLESS_SERVICE_NAME}.${env:NAMESPACE}.svc.${env:CLUSTER_DOMAIN}:${env:KAFKA_CLIENT_PORT}".to_string(),
50+
),
51+
(
52+
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP.to_string(),
53+
listener_config
54+
.listener_security_protocol_map_for_controller()),
55+
]);
56+
57+
result.insert(
58+
"inter.broker.listener.name".to_string(),
59+
KafkaListenerName::Internal.to_string(),
60+
);
61+
62+
// The ZooKeeper connection is needed for migration from ZooKeeper to KRaft mode.
63+
// It is not needed once the controller is fully running in KRaft mode.
64+
if !kraft_mode {
65+
result.insert(
66+
"zookeeper.connect".to_string(),
67+
"${env:ZOOKEEPER}".to_string(),
68+
);
69+
}
70+
71+
result.extend(kafka_security.controller_config_settings());
72+
result.extend(graceful_shutdown_config_properties());
73+
result.extend(overrides);
74+
75+
Ok(result)
76+
}
Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
//! Property-file builders for Kafka rolegroup ConfigMaps.
2+
3+
pub mod broker_properties;
4+
pub mod controller_properties;
5+
6+
use snafu::Snafu;
7+
8+
use crate::crd::{KafkaPodDescriptor, role::KafkaRole};
9+
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
18+
.iter()
19+
.filter(|pd| pd.role == KafkaRole::Controller.to_string())
20+
.map(|desc| {
21+
format!(
22+
"{fqdn}:{client_port}",
23+
fqdn = desc.fqdn(),
24+
client_port = desc.client_port
25+
)
26+
})
27+
.collect::<Vec<String>>()
28+
.join(",");
29+
30+
if result.is_empty() {
31+
None
32+
} else {
33+
Some(result)
34+
}
35+
}

0 commit comments

Comments
 (0)