Skip to content

Commit 934f6e1

Browse files
committed
zookeeper working again
1 parent f324409 commit 934f6e1

5 files changed

Lines changed: 110 additions & 39 deletions

File tree

rust/operator-binary/src/config/command.rs

Lines changed: 65 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -18,35 +18,26 @@ use crate::crd::{
1818
KAFKA_LISTENERS, KAFKA_NODE_ID, KafkaRole, broker::BROKER_PROPERTIES_FILE,
1919
controller::CONTROLLER_PROPERTIES_FILE,
2020
},
21+
v1alpha1,
2122
};
2223

2324
/// Returns the commands to start the main Kafka container
2425
pub fn broker_kafka_container_commands(
26+
kafka: &v1alpha1::KafkaCluster,
2527
cluster_id: &str,
2628
controller_descriptors: Vec<KafkaPodDescriptor>,
2729
kafka_listeners: &KafkaListenerConfig,
2830
opa_connect_string: Option<&str>,
2931
kerberos_enabled: bool,
3032
) -> String {
31-
// TODO: copy to tmp? mount readwrite folder?
3233
formatdoc! {"
3334
{COMMON_BASH_TRAP_FUNCTIONS}
3435
{remove_vector_shutdown_file_command}
3536
prepare_signal_handlers
3637
containerdebug --output={STACKABLE_LOG_DIR}/containerdebug-state.json --loop &
3738
{set_realm_env}
3839
39-
export REPLICA_ID=$(echo \"$POD_NAME\" | grep -oE '[0-9]+$')
40-
cp {config_dir}/{properties_file} /tmp/{properties_file}
41-
42-
echo \"{KAFKA_NODE_ID}=$((REPLICA_ID + {KAFKA_BROKER_ID_OFFSET}))\" >> /tmp/{properties_file}
43-
echo \"{KAFKA_CONTROLLER_QUORUM_BOOTSTRAP_SERVERS}={bootstrap_servers}\" >> /tmp/{properties_file}
44-
echo \"{KAFKA_LISTENERS}={listeners}\" >> /tmp/{properties_file}
45-
echo \"{KAFKA_ADVERTISED_LISTENERS}={advertised_listeners}\" >> /tmp/{properties_file}
46-
echo \"{KAFKA_LISTENER_SECURITY_PROTOCOL_MAP}={listener_security_protocol_map}\" >> /tmp/{properties_file}
47-
48-
bin/kafka-storage.sh format --cluster-id {cluster_id} --config /tmp/{properties_file} --initial-controllers {initial_controllers} --ignore-formatted
49-
bin/kafka-server-start.sh /tmp/{properties_file} {opa_config}{jaas_config} &
40+
{broker_start_command}
5041
5142
wait_for_termination $!
5243
{create_vector_shutdown_file_command}
@@ -57,26 +48,75 @@ pub fn broker_kafka_container_commands(
5748
true => format!("export KERBEROS_REALM=$(grep -oP 'default_realm = \\K.*' {})", STACKABLE_KERBEROS_KRB5_PATH),
5849
false => "".to_string(),
5950
},
51+
broker_start_command = broker_start_command(kafka, cluster_id, controller_descriptors, kafka_listeners, opa_connect_string, kerberos_enabled),
52+
}
53+
}
54+
55+
fn broker_start_command(
56+
kafka: &v1alpha1::KafkaCluster,
57+
cluster_id: &str,
58+
controller_descriptors: Vec<KafkaPodDescriptor>,
59+
kafka_listeners: &KafkaListenerConfig,
60+
opa_connect_string: Option<&str>,
61+
kerberos_enabled: bool,
62+
) -> String {
63+
let opa_config = match opa_connect_string {
64+
None => "".to_string(),
65+
Some(opa_connect_string) => {
66+
format!(" --override \"opa.authorizer.url={opa_connect_string}\"")
67+
}
68+
};
69+
70+
let jaas_config = match kerberos_enabled {
71+
true => {
72+
let service_name = KafkaRole::Broker.kerberos_service_name();
73+
let broker_address = node_address_cmd(STACKABLE_LISTENER_BROKER_DIR);
74+
let bootstrap_address = node_address_cmd(STACKABLE_LISTENER_BOOTSTRAP_DIR);
75+
// TODO replace client and bootstrap below with constants
76+
format!(" --override \"listener.name.client.gssapi.sasl.jaas.config=com.sun.security.auth.module.Krb5LoginModule required useKeyTab=true storeKey=true isInitiator=false keyTab=\\\"/stackable/kerberos/keytab\\\" principal=\\\"{service_name}/{broker_address}@$KERBEROS_REALM\\\";\" --override \"listener.name.bootstrap.gssapi.sasl.jaas.config=com.sun.security.auth.module.Krb5LoginModule required useKeyTab=true storeKey=true isInitiator=false keyTab=\\\"/stackable/kerberos/keytab\\\" principal=\\\"{service_name}/{bootstrap_address}@$KERBEROS_REALM\\\";\"").to_string()
77+
}
78+
false => "".to_string(),
79+
};
80+
81+
// TODO: copy to tmp? mount readwrite folder?
82+
if kafka.is_controller_configured() {
83+
formatdoc! {"
84+
export REPLICA_ID=$(echo \"$POD_NAME\" | grep -oE '[0-9]+$')
85+
cp {config_dir}/{properties_file} /tmp/{properties_file}
86+
87+
echo \"{KAFKA_NODE_ID}=$((REPLICA_ID + {KAFKA_BROKER_ID_OFFSET}))\" >> /tmp/{properties_file}
88+
echo \"{KAFKA_CONTROLLER_QUORUM_BOOTSTRAP_SERVERS}={bootstrap_servers}\" >> /tmp/{properties_file}
89+
echo \"{KAFKA_LISTENERS}={listeners}\" >> /tmp/{properties_file}
90+
echo \"{KAFKA_ADVERTISED_LISTENERS}={advertised_listeners}\" >> /tmp/{properties_file}
91+
echo \"{KAFKA_LISTENER_SECURITY_PROTOCOL_MAP}={listener_security_protocol_map}\" >> /tmp/{properties_file}
92+
93+
bin/kafka-storage.sh format --cluster-id {cluster_id} --config /tmp/{properties_file} --initial-controllers {initial_controllers} --ignore-formatted
94+
bin/kafka-server-start.sh /tmp/{properties_file} {opa_config}{jaas_config} &
95+
",
6096
config_dir = STACKABLE_CONFIG_DIR,
6197
properties_file = BROKER_PROPERTIES_FILE,
6298
bootstrap_servers = to_bootstrap_servers(&controller_descriptors),
6399
initial_controllers = to_initial_controllers(&controller_descriptors),
64100
listeners = kafka_listeners.listeners(),
65101
advertised_listeners = kafka_listeners.advertised_listeners(),
66102
listener_security_protocol_map = kafka_listeners.listener_security_protocol_map(),
67-
opa_config = match opa_connect_string {
68-
None => "".to_string(),
69-
Some(opa_connect_string) => format!(" --override \"opa.authorizer.url={opa_connect_string}\""),
70-
},
71-
jaas_config = match kerberos_enabled {
72-
true => {
73-
let service_name = KafkaRole::Broker.kerberos_service_name();
74-
let broker_address = node_address_cmd(STACKABLE_LISTENER_BROKER_DIR);
75-
let bootstrap_address = node_address_cmd(STACKABLE_LISTENER_BOOTSTRAP_DIR);
76-
// TODO replace client and bootstrap below with constants
77-
format!(" --override \"listener.name.client.gssapi.sasl.jaas.config=com.sun.security.auth.module.Krb5LoginModule required useKeyTab=true storeKey=true isInitiator=false keyTab=\\\"/stackable/kerberos/keytab\\\" principal=\\\"{service_name}/{broker_address}@$KERBEROS_REALM\\\";\" --override \"listener.name.bootstrap.gssapi.sasl.jaas.config=com.sun.security.auth.module.Krb5LoginModule required useKeyTab=true storeKey=true isInitiator=false keyTab=\\\"/stackable/kerberos/keytab\\\" principal=\\\"{service_name}/{bootstrap_address}@$KERBEROS_REALM\\\";\"").to_string()},
78-
false => "".to_string(),
79-
},
103+
}
104+
} else {
105+
formatdoc! {"
106+
bin/kafka-server-start.sh {config_dir}/{properties_file} \
107+
--override \"zookeeper.connect=$ZOOKEEPER\" \
108+
--override \"listeners={listeners}\" \
109+
--override \"advertised.listeners={advertised_listeners}\" \
110+
--override \"listener.security.protocol.map={listener_security_protocol_map}\" \
111+
{opa_config} \
112+
{jaas_config} \
113+
&",
114+
config_dir = STACKABLE_CONFIG_DIR,
115+
properties_file = BROKER_PROPERTIES_FILE,
116+
listeners = kafka_listeners.listeners(),
117+
advertised_listeners = kafka_listeners.advertised_listeners(),
118+
listener_security_protocol_map = kafka_listeners.listener_security_protocol_map(),
119+
}
80120
}
81121
}
82122

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -235,6 +235,7 @@ pub fn get_kafka_listener_config(
235235

236236
// CONTROLLER
237237
if kafka.is_controller_configured() {
238+
// TODO: SSL?
238239
listener_security_protocol_map.insert(
239240
KafkaListenerName::Controller,
240241
KafkaListenerProtocol::Plaintext,

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

Lines changed: 11 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -137,21 +137,23 @@ impl Configuration for BrokerConfigFragment {
137137
let mut config = BTreeMap::new();
138138

139139
if file == BROKER_PROPERTIES_FILE {
140-
config.insert(
141-
KAFKA_PROCESS_ROLES.to_string(),
142-
Some(KafkaRole::Broker.to_string()),
143-
);
144-
145140
config.insert(
146141
KAFKA_LOG_DIRS.to_string(),
147142
Some("/stackable/data/topicdata".to_string()),
148143
);
149144

150-
config.insert(
151-
"controller.listener.names".to_string(),
152-
Some(KafkaListenerName::Controller.to_string()),
153-
);
145+
// KRAFT
146+
if resource.is_controller_configured() {
147+
config.insert(
148+
KAFKA_PROCESS_ROLES.to_string(),
149+
Some(KafkaRole::Broker.to_string()),
150+
);
154151

152+
config.insert(
153+
"controller.listener.names".to_string(),
154+
Some(KafkaListenerName::Controller.to_string()),
155+
);
156+
}
155157
// OPA
156158
if resource.spec.cluster_config.authorization.opa.is_some() {
157159
config.insert(

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

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ use stackable_operator::{
1515
use strum::{Display, EnumIter};
1616

1717
use crate::crd::{
18+
listener::KafkaListenerName,
1819
role::{
1920
KAFKA_LOG_DIRS, KAFKA_PROCESS_ROLES, KafkaRole,
2021
commons::{CommonConfig, Storage, StorageFragment},
@@ -127,14 +128,20 @@ impl Configuration for ControllerConfigFragment {
127128
let mut config = BTreeMap::new();
128129

129130
if file == CONTROLLER_PROPERTIES_FILE {
131+
config.insert(
132+
KAFKA_LOG_DIRS.to_string(),
133+
Some("/stackable/data/kraft".to_string()),
134+
);
135+
136+
// KRAFT
130137
config.insert(
131138
KAFKA_PROCESS_ROLES.to_string(),
132139
Some(KafkaRole::Controller.to_string()),
133140
);
134141

135142
config.insert(
136-
KAFKA_LOG_DIRS.to_string(),
137-
Some("/stackable/data/kraft".to_string()),
143+
"controller.listener.names".to_string(),
144+
Some(KafkaListenerName::Controller.to_string()),
138145
);
139146

140147
config.insert(

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

Lines changed: 24 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -23,10 +23,11 @@ use stackable_operator::{
2323
apps::v1::{StatefulSet, StatefulSetSpec},
2424
core::v1::{
2525
ConfigMapKeySelector, ConfigMapVolumeSource, ContainerPort, EnvVar, EnvVarSource,
26-
ExecAction, ObjectFieldSelector, PodSpec, Probe, ServiceAccount, Volume,
26+
ExecAction, ObjectFieldSelector, PodSpec, Probe, ServiceAccount, TCPSocketAction,
27+
Volume,
2728
},
2829
},
29-
apimachinery::pkg::apis::meta::v1::LabelSelector,
30+
apimachinery::pkg::{apis::meta::v1::LabelSelector, util::intstr::IntOrString},
3031
},
3132
kube::ResourceExt,
3233
kvp::Labels,
@@ -292,6 +293,7 @@ pub fn build_broker_rolegroup_statefulset(
292293
"-c".to_string(),
293294
])
294295
.args(vec![broker_kafka_container_commands(
296+
kafka,
295297
cluster_id,
296298
// we need controller pods
297299
kafka
@@ -661,7 +663,26 @@ pub fn build_controller_rolegroup_statefulset(
661663
.context(AddVolumeMountSnafu)?
662664
.add_volume_mount("log", STACKABLE_LOG_DIR)
663665
.context(AddVolumeMountSnafu)?
664-
.resources(merged_config.resources().clone().into());
666+
.resources(merged_config.resources().clone().into())
667+
// TODO: improve probes
668+
.liveness_probe(Probe {
669+
tcp_socket: Some(TCPSocketAction {
670+
port: IntOrString::Int(kafka_security.client_port().into()),
671+
..Default::default()
672+
}),
673+
timeout_seconds: Some(5),
674+
period_seconds: Some(5),
675+
..Probe::default()
676+
})
677+
.readiness_probe(Probe {
678+
tcp_socket: Some(TCPSocketAction {
679+
port: IntOrString::Int(kafka_security.client_port().into()),
680+
..Default::default()
681+
}),
682+
timeout_seconds: Some(5),
683+
period_seconds: Some(5),
684+
..Probe::default()
685+
});
665686

666687
if let ContainerLogConfig {
667688
choice:

0 commit comments

Comments
 (0)