Skip to content

Commit c49b949

Browse files
committed
remove wrongly added pod names to discovery map bootstrap hosts
1 parent 3af62b9 commit c49b949

2 files changed

Lines changed: 50 additions & 33 deletions

File tree

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

Lines changed: 40 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -189,6 +189,15 @@ impl Display for KafkaListener {
189189
}
190190
}
191191

192+
// Builds a list of listeners for the given Kafka cluster and rolegroup.
193+
//
194+
// TODO: Not every listener is necessarily used by every role while some listeners are used by both roles.
195+
// Yeah, this is confusing and needs refactoring.
196+
//
197+
// For example, the INTERNAL and CLIENT listener are only configured on brokers.
198+
// On the other hand, the BOOTSTRAP listener is configured on both roles.
199+
// Note that actual the bootstrap services are different between brokers and controllers but from the
200+
// Kafka perspective they are both called BOOTSTRAP.
192201
pub fn get_kafka_listener_config(
193202
kafka: &v1alpha1::KafkaCluster,
194203
kafka_security: &KafkaTlsSecurity,
@@ -301,7 +310,7 @@ pub fn get_kafka_listener_config(
301310
);
302311
}
303312

304-
let bootstrap_pod_fqdn = pod_fqdn(
313+
let bootstrap_service_fqdn = service_fqdn(
305314
kafka,
306315
&rolegroup_service_name(rolegroup_ref, KafkaListenerName::Bootstrap),
307316
cluster_info,
@@ -315,7 +324,7 @@ pub fn get_kafka_listener_config(
315324
});
316325
advertised_listeners.push(KafkaListener {
317326
name: KafkaListenerName::Bootstrap,
318-
host: bootstrap_pod_fqdn.to_string(),
327+
host: bootstrap_service_fqdn.to_string(),
319328
port: node_port_cmd(
320329
STACKABLE_LISTENER_BOOTSTRAP_DIR,
321330
kafka_security.bootstrap_port_name(),
@@ -343,19 +352,6 @@ pub fn get_kafka_listener_config(
343352
})
344353
}
345354

346-
// TODO: This is the more general version to `RoleGroupRef::rolegroup_headless_service_name()`
347-
// because we need it for the bootstrap service as well, which doesn't exist in op-rs.
348-
fn rolegroup_service_name(
349-
rolegroup_ref: &RoleGroupRef<v1alpha1::KafkaCluster>,
350-
listener: KafkaListenerName,
351-
) -> String {
352-
format!(
353-
"{name}-{service}",
354-
name = rolegroup_ref.object_name(),
355-
service = listener.to_string().to_lowercase()
356-
)
357-
}
358-
359355
pub fn node_address_cmd_env(directory: &str) -> String {
360356
format!("$(cat {directory}/default-address/address)")
361357
}
@@ -384,6 +380,31 @@ pub fn pod_fqdn(
384380
))
385381
}
386382

383+
// TODO: This is the more general version to `RoleGroupRef::rolegroup_headless_service_name()`
384+
// because we need it for the bootstrap service as well, which doesn't exist in op-rs.
385+
fn rolegroup_service_name(
386+
rolegroup_ref: &RoleGroupRef<v1alpha1::KafkaCluster>,
387+
listener: KafkaListenerName,
388+
) -> String {
389+
format!(
390+
"{name}-{service}",
391+
name = rolegroup_ref.object_name(),
392+
service = listener.to_string().to_lowercase()
393+
)
394+
}
395+
396+
pub fn service_fqdn(
397+
kafka: &v1alpha1::KafkaCluster,
398+
sts_service_name: &str,
399+
cluster_info: &KubernetesClusterInfo,
400+
) -> Result<String, KafkaListenerError> {
401+
Ok(format!(
402+
"{sts_service_name}.{namespace}.svc.{cluster_domain}",
403+
namespace = kafka.namespace().context(ObjectHasNoNamespaceSnafu)?,
404+
cluster_domain = cluster_info.cluster_domain
405+
))
406+
}
407+
387408
#[cfg(test)]
388409
mod tests {
389410
use stackable_operator::{
@@ -478,7 +499,7 @@ mod tests {
478499
.unwrap(),
479500
internal_port = kafka_security.internal_port(),
480501
bootstrap_name = KafkaListenerName::Bootstrap,
481-
bootstrap_host = pod_fqdn(
502+
bootstrap_host = service_fqdn(
482503
&kafka,
483504
&rolegroup_service_name(&rolegroup_ref, KafkaListenerName::Bootstrap),
484505
&cluster_info
@@ -550,7 +571,7 @@ mod tests {
550571
.unwrap(),
551572
internal_port = kafka_security.internal_port(),
552573
bootstrap_name = KafkaListenerName::Bootstrap,
553-
bootstrap_host = pod_fqdn(
574+
bootstrap_host = service_fqdn(
554575
&kafka,
555576
&rolegroup_service_name(&rolegroup_ref, KafkaListenerName::Bootstrap),
556577
&cluster_info
@@ -623,7 +644,7 @@ mod tests {
623644
.unwrap(),
624645
internal_port = kafka_security.internal_port(),
625646
bootstrap_name = KafkaListenerName::Bootstrap,
626-
bootstrap_host = pod_fqdn(
647+
bootstrap_host = service_fqdn(
627648
&kafka,
628649
&rolegroup_service_name(&rolegroup_ref, KafkaListenerName::Bootstrap),
629650
&cluster_info
@@ -728,7 +749,7 @@ mod tests {
728749
.unwrap(),
729750
internal_port = kafka_security.internal_port(),
730751
bootstrap_name = KafkaListenerName::Bootstrap,
731-
bootstrap_host = pod_fqdn(
752+
bootstrap_host = service_fqdn(
732753
&kafka,
733754
&rolegroup_service_name(&rolegroup_ref, KafkaListenerName::Bootstrap),
734755
&cluster_info

rust/operator-binary/src/discovery.rs

Lines changed: 10 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -52,19 +52,15 @@ pub fn build_discovery_configmap(
5252
kafka_security: &KafkaTlsSecurity,
5353
listeners: &[listener::v1alpha1::Listener],
5454
) -> Result<ConfigMap, Error> {
55-
let port_name = if kafka_security.has_kerberos_enabled() {
56-
kafka_security.bootstrap_port_name()
57-
} else {
58-
kafka_security.client_port_name()
59-
};
60-
61-
// Write a list of bootstrap servers in the format that Kafka clients:
62-
// "{host1}:{port1},{host2:port2},..."
63-
let bootstrap_servers = listener_hosts(listeners, port_name)?
64-
.into_iter()
65-
.map(|(host, port)| format!("{}:{}", host, port))
66-
.collect::<Vec<_>>()
67-
.join(",");
55+
// Write a list of *broker* bootstrap hosts for all rolegroups separated by commas:
56+
// "{host1}:{port},{host2:port},..."
57+
// The port is the same bootstrap port for all services.
58+
let bootstrap_servers =
59+
listener_hosts_filtered_by_port_name(listeners, kafka_security.bootstrap_port_name())?
60+
.into_iter()
61+
.map(|(host, port)| format!("{}:{}", host, port))
62+
.collect::<Vec<_>>()
63+
.join(",");
6864
ConfigMapBuilder::new()
6965
.metadata(
7066
ObjectMetaBuilder::new()
@@ -89,7 +85,7 @@ pub fn build_discovery_configmap(
8985
.context(BuildConfigMapSnafu)
9086
}
9187

92-
fn listener_hosts(
88+
fn listener_hosts_filtered_by_port_name(
9389
listeners: &[listener::v1alpha1::Listener],
9490
port_name: &str,
9591
) -> Result<impl IntoIterator<Item = (String, u16)>, Error> {

0 commit comments

Comments
 (0)