Skip to content

Commit 1b7a6a4

Browse files
committed
add headleass & metrics service
1 parent e2ac187 commit 1b7a6a4

4 files changed

Lines changed: 12 additions & 21 deletions

File tree

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -187,7 +187,7 @@ pub fn controller_kafka_container_command(
187187
fn to_listeners(port: u16) -> String {
188188
// The environment variables are set in the statefulset of the controller
189189
format!(
190-
"{listener_name}://$POD_NAME.$ROLEGROUP_REF.$NAMESPACE.svc.$CLUSTER_DOMAIN:{port}",
190+
"{listener_name}://$POD_NAME.$ROLEGROUP_HEADLESS_SERVICE_NAME.$NAMESPACE.svc.$CLUSTER_DOMAIN:{port}",
191191
listener_name = KafkaListenerName::Controller
192192
)
193193
}

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

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -294,7 +294,9 @@ impl v1alpha1::KafkaCluster {
294294
for replica in 0..replicas {
295295
pod_descriptors.push(KafkaPodDescriptor {
296296
namespace: namespace.clone(),
297-
role_group_service_name: rolegroup_ref.object_name(),
297+
role_group_service_name: rolegroup_ref
298+
.rolegroup_headless_service_name(),
299+
role_group_statefulset_name: rolegroup_ref.object_name(),
298300
replica,
299301
cluster_domain: cluster_info.cluster_domain.clone(),
300302
node_id: node_id_hash_offset + u32::from(replica),
@@ -341,6 +343,7 @@ impl v1alpha1::KafkaCluster {
341343
#[derive(Debug, PartialEq, Eq, PartialOrd, Ord)]
342344
pub struct KafkaPodDescriptor {
343345
namespace: String,
346+
role_group_statefulset_name: String,
344347
role_group_service_name: String,
345348
replica: u16,
346349
cluster_domain: DomainName,
@@ -361,7 +364,7 @@ impl KafkaPodDescriptor {
361364
}
362365

363366
pub fn pod_name(&self) -> String {
364-
format!("{}-{}", self.role_group_service_name, self.replica)
367+
format!("{}-{}", self.role_group_statefulset_name, self.replica)
365368
}
366369

367370
/// Build the Kraft voter String

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

Lines changed: 2 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,7 @@ pub fn build_rolegroup_headless_service(
4343
Ok(Service {
4444
metadata: ObjectMetaBuilder::new()
4545
.name_and_namespace(kafka)
46-
.name(rolegroup_headless_service_name(rolegroup))
46+
.name(rolegroup.rolegroup_headless_service_name())
4747
.ownerreference_from_resource(kafka, None, Some(true))
4848
.context(ObjectMissingMetadataForOwnerRefSnafu)?
4949
.with_recommended_labels(build_recommended_labels(
@@ -85,7 +85,7 @@ pub fn build_rolegroup_metrics_service(
8585
metadata: ObjectMetaBuilder::new()
8686
.name_and_namespace(kafka)
8787
// TODO: Use method on RoleGroupRef once op-rs is released
88-
.name(rolegroup_metrics_service_name(rolegroup))
88+
.name(rolegroup.rolegroup_metrics_service_name())
8989
.ownerreference_from_resource(kafka, None, Some(true))
9090
.context(ObjectMissingMetadataForOwnerRefSnafu)?
9191
.with_recommended_labels(build_recommended_labels(
@@ -122,18 +122,6 @@ pub fn build_rolegroup_metrics_service(
122122
Ok(metrics_service)
123123
}
124124

125-
/// Headless service for cluster internal purposes only.
126-
// TODO: Move to operator-rs
127-
fn rolegroup_headless_service_name(rolegroup: &RoleGroupRef<v1alpha1::KafkaCluster>) -> String {
128-
format!("{name}-headless", name = rolegroup.object_name())
129-
}
130-
131-
/// Headless metrics service exposes Prometheus endpoint only
132-
// TODO: Move to operator-rs
133-
fn rolegroup_metrics_service_name(rolegroup: &RoleGroupRef<v1alpha1::KafkaCluster>) -> String {
134-
format!("{name}-metrics", name = rolegroup.object_name())
135-
}
136-
137125
fn metrics_ports() -> Vec<ServicePort> {
138126
vec![ServicePort {
139127
name: Some(METRICS_PORT_NAME.to_string()),

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

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -547,7 +547,7 @@ pub fn build_broker_rolegroup_statefulset(
547547
),
548548
..LabelSelector::default()
549549
},
550-
service_name: Some(rolegroup_ref.object_name()),
550+
service_name: Some(rolegroup_ref.rolegroup_headless_service_name()),
551551
template: pod_template,
552552
volume_claim_templates: Some(pvcs),
553553
..StatefulSetSpec::default()
@@ -621,8 +621,8 @@ pub fn build_controller_rolegroup_statefulset(
621621
});
622622

623623
env.push(EnvVar {
624-
name: "ROLEGROUP_REF".to_string(),
625-
value: Some(rolegroup_ref.object_name()),
624+
name: "ROLEGROUP_HEADLESS_SERVICE_NAME".to_string(),
625+
value: Some(rolegroup_ref.rolegroup_headless_service_name()),
626626
..EnvVar::default()
627627
});
628628

@@ -875,7 +875,7 @@ pub fn build_controller_rolegroup_statefulset(
875875
),
876876
..LabelSelector::default()
877877
},
878-
service_name: Some(rolegroup_ref.object_name()),
878+
service_name: Some(rolegroup_ref.rolegroup_headless_service_name()),
879879
template: pod_template,
880880
volume_claim_templates: Some(merged_config.resources().storage.build_pvcs()),
881881
..StatefulSetSpec::default()

0 commit comments

Comments
 (0)