Skip to content

Commit f33fbc1

Browse files
committed
wip - controller tls
1 parent 104c1b1 commit f33fbc1

13 files changed

Lines changed: 313 additions & 124 deletions

File tree

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

Lines changed: 17 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@ use stackable_operator::{
99
use crate::crd::{
1010
KafkaPodDescriptor, STACKABLE_CONFIG_DIR, STACKABLE_KERBEROS_KRB5_PATH,
1111
STACKABLE_LISTENER_BOOTSTRAP_DIR, STACKABLE_LISTENER_BROKER_DIR, STACKABLE_LOG_DIR,
12-
listener::{KafkaListenerConfig, node_address_cmd},
12+
listener::{KafkaListenerConfig, KafkaListenerName, node_address_cmd},
1313
role::{
1414
KAFKA_ADVERTISED_LISTENERS, KAFKA_BROKER_ID_OFFSET,
1515
KAFKA_CONTROLLER_QUORUM_BOOTSTRAP_SERVERS, KAFKA_LISTENER_SECURITY_PROTOCOL_MAP,
@@ -106,9 +106,9 @@ fn broker_start_command(
106106
formatdoc! {"
107107
bin/kafka-server-start.sh {config_dir}/{properties_file} \
108108
--override \"zookeeper.connect=$ZOOKEEPER\" \
109-
--override \"listeners={listeners}\" \
110-
--override \"advertised.listeners={advertised_listeners}\" \
111-
--override \"listener.security.protocol.map={listener_security_protocol_map}\" \
109+
--override \"{KAFKA_LISTENERS}={listeners}\" \
110+
--override \"{KAFKA_ADVERTISED_LISTENERS}={advertised_listeners}\" \
111+
--override \"{KAFKA_LISTENER_SECURITY_PROTOCOL_MAP}={listener_security_protocol_map}\" \
112112
{opa_config} \
113113
{jaas_config} \
114114
&",
@@ -124,6 +124,7 @@ fn broker_start_command(
124124
pub fn controller_kafka_container_command(
125125
cluster_id: &str,
126126
controller_descriptors: Vec<KafkaPodDescriptor>,
127+
kafka_listeners: &KafkaListenerConfig,
127128
kafka_security: &KafkaTlsSecurity,
128129
) -> String {
129130
let client_port = kafka_security.client_port();
@@ -153,7 +154,7 @@ pub fn controller_kafka_container_command(
153154
properties_file = CONTROLLER_PROPERTIES_FILE,
154155
bootstrap_servers = to_bootstrap_servers(&controller_descriptors, client_port),
155156
listeners = to_listeners(client_port),
156-
listener_security_protocol_map = to_listener_security_protocol_map(),
157+
listener_security_protocol_map = to_listener_security_protocol_map(kafka_listeners),
157158
initial_controllers = to_initial_controllers(&controller_descriptors, client_port),
158159
create_vector_shutdown_file_command = create_vector_shutdown_file_command(STACKABLE_LOG_DIR)
159160
}
@@ -162,13 +163,19 @@ pub fn controller_kafka_container_command(
162163
fn to_listeners(port: u16) -> String {
163164
// TODO:
164165
// - document that variables are set in stateful set
165-
// - customize listener (CONTROLLER)
166-
format!("CONTROLLER://$POD_NAME.$ROLEGROUP_REF.$NAMESPACE.svc.$CLUSTER_DOMAIN:{port}")
166+
// - customize listener (CONTROLLER / CONTROLLER_AUTH?)
167+
format!(
168+
"{listener_name}://$POD_NAME.$ROLEGROUP_REF.$NAMESPACE.svc.$CLUSTER_DOMAIN:{port}",
169+
listener_name = KafkaListenerName::Controller
170+
)
167171
}
168172

169-
fn to_listener_security_protocol_map() -> String {
170-
// TODO: make configurable
171-
"CONTROLLER:PLAINTEXT".to_string()
173+
fn to_listener_security_protocol_map(kafka_listeners: &KafkaListenerConfig) -> String {
174+
// TODO: make configurable - CONTROLLER_AUTH
175+
kafka_listeners
176+
.listener_security_protocol_map_for_listener(&KafkaListenerName::Controller)
177+
// todo better error
178+
.unwrap_or("".to_string())
172179
}
173180

174181
fn to_initial_controllers(controller_descriptors: &[KafkaPodDescriptor], port: u16) -> String {
@@ -186,11 +193,3 @@ fn to_bootstrap_servers(controller_descriptors: &[KafkaPodDescriptor], port: u16
186193
.collect::<Vec<String>>()
187194
.join(",")
188195
}
189-
190-
// fn to_kafka_overrides(overrides: BTreeMap<String, String>) -> String {
191-
// overrides
192-
// .iter()
193-
// .map(|(key, value)| format!("--override \"{key}={value}\""))
194-
// .collect::<Vec<String>>()
195-
// .join(" ")
196-
// }

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

Lines changed: 78 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,59 @@ pub enum KafkaListenerName {
4444
Bootstrap,
4545
#[strum(serialize = "CONTROLLER")]
4646
Controller,
47+
#[strum(serialize = "CONTROLLER_AUTH")]
48+
ControllerAuth,
49+
}
50+
51+
impl KafkaListenerName {
52+
pub fn listener_ssl_keystore_location(&self) -> String {
53+
format!(
54+
"listener.name.{listener_name}.ssl.keystore.location",
55+
listener_name = self.to_string().to_lowercase()
56+
)
57+
}
58+
59+
pub fn listener_ssl_keystore_password(&self) -> String {
60+
format!(
61+
"listener.name.{listener_name}.ssl.keystore.password",
62+
listener_name = self.to_string().to_lowercase()
63+
)
64+
}
65+
66+
pub fn listener_ssl_keystore_type(&self) -> String {
67+
format!(
68+
"listener.name.{listener_name}.ssl.keystore.type",
69+
listener_name = self.to_string().to_lowercase()
70+
)
71+
}
72+
73+
pub fn listener_ssl_truststore_location(&self) -> String {
74+
format!(
75+
"listener.name.{listener_name}.ssl.truststore.location",
76+
listener_name = self.to_string().to_lowercase()
77+
)
78+
}
79+
80+
pub fn listener_ssl_truststore_password(&self) -> String {
81+
format!(
82+
"listener.name.{listener_name}.ssl.truststore.password",
83+
listener_name = self.to_string().to_lowercase()
84+
)
85+
}
86+
87+
pub fn listener_ssl_truststore_type(&self) -> String {
88+
format!(
89+
"listener.name.{listener_name}.ssl.truststore.type",
90+
listener_name = self.to_string().to_lowercase()
91+
)
92+
}
93+
94+
pub fn listener_ssl_client_auth(&self) -> String {
95+
format!(
96+
"listener.name.{listener_name}.ssl.client.auth",
97+
listener_name = self.to_string().to_lowercase()
98+
)
99+
}
47100
}
48101

49102
#[derive(Debug)]
@@ -82,6 +135,16 @@ impl KafkaListenerConfig {
82135
.collect::<Vec<String>>()
83136
.join(",")
84137
}
138+
139+
/// Returns the `listener.security.protocol.map` for the Kafka `broker.properties` config.
140+
pub fn listener_security_protocol_map_for_listener(
141+
&self,
142+
listener_name: &KafkaListenerName,
143+
) -> Option<String> {
144+
self.listener_security_protocol_map
145+
.get(listener_name)
146+
.map(|protocol| format!("{listener_name}:{protocol}"))
147+
}
85148
}
86149

87150
#[derive(Debug)]
@@ -127,6 +190,10 @@ pub fn get_kafka_listener_config(
127190
});
128191
listener_security_protocol_map
129192
.insert(KafkaListenerName::ClientAuth, KafkaListenerProtocol::Ssl);
193+
listener_security_protocol_map.insert(
194+
KafkaListenerName::ControllerAuth,
195+
KafkaListenerProtocol::Ssl,
196+
);
130197
} else if kafka_security.has_kerberos_enabled() {
131198
// 2) Kerberos and TLS authentication classes are mutually exclusive
132199
listeners.push(KafkaListener {
@@ -144,6 +211,10 @@ pub fn get_kafka_listener_config(
144211
});
145212
listener_security_protocol_map
146213
.insert(KafkaListenerName::Client, KafkaListenerProtocol::SaslSsl);
214+
listener_security_protocol_map.insert(
215+
KafkaListenerName::Controller,
216+
KafkaListenerProtocol::SaslSsl,
217+
);
147218
} else if kafka_security.tls_server_secret_class().is_some() {
148219
// 3) If no client authentication but tls is required we expose CLIENT with SSL
149220
listeners.push(KafkaListener {
@@ -180,7 +251,7 @@ pub fn get_kafka_listener_config(
180251
.insert(KafkaListenerName::Client, KafkaListenerProtocol::Plaintext);
181252
}
182253

183-
// INTERNAL
254+
// INTERNAL / CONTROLLER
184255
if kafka_security.has_kerberos_enabled() || kafka_security.tls_internal_secret_class().is_some()
185256
{
186257
// 5) & 6) Kerberos and TLS authentication classes are mutually exclusive but both require internal tls to be used
@@ -196,6 +267,8 @@ pub fn get_kafka_listener_config(
196267
});
197268
listener_security_protocol_map
198269
.insert(KafkaListenerName::Internal, KafkaListenerProtocol::Ssl);
270+
listener_security_protocol_map
271+
.insert(KafkaListenerName::Controller, KafkaListenerProtocol::Ssl);
199272
} else {
200273
// 7) If no internal tls is required we expose INTERNAL as PLAINTEXT
201274
listeners.push(KafkaListener {
@@ -212,6 +285,10 @@ pub fn get_kafka_listener_config(
212285
KafkaListenerName::Internal,
213286
KafkaListenerProtocol::Plaintext,
214287
);
288+
listener_security_protocol_map.insert(
289+
KafkaListenerName::Controller,
290+
KafkaListenerProtocol::Plaintext,
291+
);
215292
}
216293

217294
// BOOTSTRAP
@@ -233,15 +310,6 @@ pub fn get_kafka_listener_config(
233310
.insert(KafkaListenerName::Bootstrap, KafkaListenerProtocol::SaslSsl);
234311
}
235312

236-
// CONTROLLER
237-
if kafka.is_controller_configured() {
238-
// TODO: SSL?
239-
listener_security_protocol_map.insert(
240-
KafkaListenerName::Controller,
241-
KafkaListenerProtocol::Plaintext,
242-
);
243-
}
244-
245313
Ok(KafkaListenerConfig {
246314
listeners,
247315
advertised_listeners,

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

Lines changed: 0 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -143,11 +143,6 @@ impl Configuration for ControllerConfigFragment {
143143
"controller.listener.names".to_string(),
144144
Some(KafkaListenerName::Controller.to_string()),
145145
);
146-
147-
config.insert(
148-
"controller.listener.names".to_string(),
149-
Some("CONTROLLER".to_string()),
150-
);
151146
}
152147

153148
Ok(config)

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -130,6 +130,7 @@ impl KafkaRole {
130130
/// A Kerberos principal has three parts, with the form username/fully.qualified.domain.name@YOUR-REALM.COM.
131131
/// We only have one role and will use "kafka" everywhere (which e.g. differs from the current hdfs implementation,
132132
/// but is similar to HBase).
133+
// TODO: split into broker / controller?
133134
pub fn kerberos_service_name(&self) -> &'static str {
134135
"kafka"
135136
}

0 commit comments

Comments
 (0)