Skip to content

Commit 9f2a57f

Browse files
committed
Moving services to own module, splitting up functions
1 parent 4298a8f commit 9f2a57f

4 files changed

Lines changed: 176 additions & 138 deletions

File tree

rust/operator-binary/src/controller.rs

Lines changed: 33 additions & 100 deletions
Original file line numberDiff line numberDiff line change
@@ -43,8 +43,8 @@ use stackable_operator::{
4343
api::{
4444
apps::v1::{StatefulSet, StatefulSetSpec},
4545
core::v1::{
46-
ConfigMap, ConfigMapVolumeSource, EmptyDirVolumeSource, Probe, Service,
47-
ServiceSpec, TCPSocketAction, Volume,
46+
ConfigMap, ConfigMapVolumeSource, EmptyDirVolumeSource, Probe, TCPSocketAction,
47+
Volume,
4848
},
4949
},
5050
apimachinery::pkg::{
@@ -56,7 +56,7 @@ use stackable_operator::{
5656
core::{DeserializeGuard, error_boundary},
5757
runtime::controller::Action,
5858
},
59-
kvp::{Label, Labels, ObjectLabels},
59+
kvp::{Labels, ObjectLabels},
6060
logging::controller::ReconcilerError,
6161
memory::{BinaryMultiple, MemoryQuantity},
6262
product_config_utils::{transform_all_roles_to_config, validate_all_roles_and_groups_config},
@@ -102,6 +102,10 @@ use crate::{
102102
listener::{LISTENER_VOLUME_DIR, LISTENER_VOLUME_NAME, build_role_listener},
103103
operations::{graceful_shutdown::add_graceful_shutdown_config, pdb::add_pdbs},
104104
product_logging::extend_role_group_config_map,
105+
service::{
106+
build_rolegroup_headless_service, build_rolegroup_metrics_service,
107+
rolegroup_metrics_service_name,
108+
},
105109
};
106110

107111
pub const HIVE_CONTROLLER_NAME: &str = "hivecluster";
@@ -345,6 +349,9 @@ pub enum Error {
345349
BuildListenerVolume {
346350
source: ListenerOperatorVolumeSourceBuilderError,
347351
},
352+
353+
#[snafu(display("faild to configure service"))]
354+
ServiceConfiguration { source: crate::service::Error },
348355
}
349356
type Result<T, E = Error> = std::result::Result<T, E>;
350357

@@ -456,7 +463,14 @@ pub async fn reconcile_hive(
456463
.merged_config(&HiveRole::MetaStore, &rolegroup)
457464
.context(FailedToResolveResourceConfigSnafu)?;
458465

459-
let rg_services = build_rolegroup_services(hive, &resolved_product_image, &rolegroup)?;
466+
let rg_metrics_service =
467+
build_rolegroup_metrics_service(hive, &resolved_product_image, &rolegroup)
468+
.context(ServiceConfigurationSnafu)?;
469+
470+
let rg_headless_service =
471+
build_rolegroup_headless_service(hive, &resolved_product_image, &rolegroup)
472+
.context(ServiceConfigurationSnafu)?;
473+
460474
let rg_configmap = build_metastore_rolegroup_config_map(
461475
hive,
462476
&hive_namespace,
@@ -478,13 +492,19 @@ pub async fn reconcile_hive(
478492
&rbac_sa.name_any(),
479493
)?;
480494

481-
for rg_service in rg_services {
482-
cluster_resources.add(client, rg_service).await.context(
483-
ApplyRoleGroupServiceSnafu {
484-
rolegroup: rolegroup.clone(),
485-
},
486-
)?;
487-
}
495+
cluster_resources
496+
.add(client, rg_metrics_service)
497+
.await
498+
.context(ApplyRoleGroupServiceSnafu {
499+
rolegroup: rolegroup.clone(),
500+
})?;
501+
502+
cluster_resources
503+
.add(client, rg_headless_service)
504+
.await
505+
.context(ApplyRoleGroupServiceSnafu {
506+
rolegroup: rolegroup.clone(),
507+
})?;
488508

489509
cluster_resources
490510
.add(client, rg_configmap)
@@ -714,97 +734,10 @@ fn build_metastore_rolegroup_config_map(
714734
})
715735
}
716736

717-
/// The rolegroup [`Service`] is a headless service that allows direct access to the instances of a certain rolegroup
718-
///
719-
/// This is mostly useful for internal communication between peers, or for clients that perform client-side load balancing.
720-
fn build_rolegroup_services(
721-
hive: &v1alpha1::HiveCluster,
722-
resolved_product_image: &ResolvedProductImage,
723-
rolegroup: &RoleGroupRef<v1alpha1::HiveCluster>,
724-
) -> Result<Vec<Service>> {
725-
let services = vec![
726-
Service {
727-
metadata: ObjectMetaBuilder::new()
728-
.name_and_namespace(hive)
729-
// TODO: Use method on RoleGroupRef once op-rs is released
730-
.name(hive.rolegroup_headless_metrics_service_name(rolegroup))
731-
.ownerreference_from_resource(hive, None, Some(true))
732-
.context(ObjectMissingMetadataForOwnerRefSnafu)?
733-
.with_recommended_labels(build_recommended_labels(
734-
hive,
735-
&resolved_product_image.app_version_label,
736-
&rolegroup.role,
737-
&rolegroup.role_group,
738-
))
739-
.context(MetadataBuildSnafu)?
740-
.with_label(
741-
Label::try_from(("prometheus.io/scrape", "true")).context(LabelBuildSnafu)?,
742-
)
743-
.build(),
744-
spec: Some(ServiceSpec {
745-
// Internal communication does not need to be exposed
746-
type_: Some("ClusterIP".to_string()),
747-
cluster_ip: Some("None".to_string()),
748-
ports: Some(hive.metrics_ports()),
749-
selector: Some(
750-
Labels::role_group_selector(
751-
hive,
752-
APP_NAME,
753-
&rolegroup.role,
754-
&rolegroup.role_group,
755-
)
756-
.context(LabelBuildSnafu)?
757-
.into(),
758-
),
759-
publish_not_ready_addresses: Some(true),
760-
..ServiceSpec::default()
761-
}),
762-
status: None,
763-
},
764-
Service {
765-
metadata: ObjectMetaBuilder::new()
766-
.name_and_namespace(hive)
767-
// TODO: Use method on RoleGroupRef once op-rs is released
768-
.name(hive.rolegroup_headless_service_name(rolegroup))
769-
.ownerreference_from_resource(hive, None, Some(true))
770-
.context(ObjectMissingMetadataForOwnerRefSnafu)?
771-
.with_recommended_labels(build_recommended_labels(
772-
hive,
773-
&resolved_product_image.app_version_label,
774-
&rolegroup.role,
775-
&rolegroup.role_group,
776-
))
777-
.context(MetadataBuildSnafu)?
778-
.build(),
779-
spec: Some(ServiceSpec {
780-
// Internal communication does not need to be exposed
781-
type_: Some("ClusterIP".to_string()),
782-
cluster_ip: Some("None".to_string()),
783-
// Expecting same ports as on listener service, just as a headless, internal service
784-
ports: Some(hive.service_ports()),
785-
selector: Some(
786-
Labels::role_group_selector(
787-
hive,
788-
APP_NAME,
789-
&rolegroup.role,
790-
&rolegroup.role_group,
791-
)
792-
.context(LabelBuildSnafu)?
793-
.into(),
794-
),
795-
publish_not_ready_addresses: Some(true),
796-
..ServiceSpec::default()
797-
}),
798-
status: None,
799-
},
800-
];
801-
Ok(services)
802-
}
803-
804737
/// The rolegroup [`StatefulSet`] runs the rolegroup, as configured by the administrator.
805738
///
806739
/// The [`Pod`](`stackable_operator::k8s_openapi::api::core::v1::Pod`)s are accessible through the
807-
/// corresponding [`Service`] (from [`build_rolegroup_services`]).
740+
/// corresponding [`Service`] (from [`build_rolegroup_headless_service`] and [`build_rolegroup_metrics_service`]).
808741
#[allow(clippy::too_many_arguments)]
809742
fn build_metastore_rolegroup_statefulset(
810743
hive: &v1alpha1::HiveCluster,
@@ -1147,7 +1080,7 @@ fn build_metastore_rolegroup_statefulset(
11471080
..LabelSelector::default()
11481081
},
11491082
// TODO: Use method on RoleGroupRef once op-rs is released
1150-
service_name: Some(hive.rolegroup_headless_metrics_service_name(rolegroup_ref)),
1083+
service_name: Some(rolegroup_metrics_service_name(rolegroup_ref)),
11511084
template: pod_template,
11521085
volume_claim_templates: Some(vec![pvc]),
11531086
..StatefulSetSpec::default()

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

Lines changed: 1 addition & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@ use stackable_operator::{
1818
merge::Merge,
1919
},
2020
crd::s3,
21-
k8s_openapi::{api::core::v1::ServicePort, apimachinery::pkg::api::resource::Quantity},
21+
k8s_openapi::apimachinery::pkg::api::resource::Quantity,
2222
kube::{CustomResource, ResourceExt, runtime::reflector::ObjectRef},
2323
product_config_utils::{self, Configuration},
2424
product_logging::{self, spec::Logging},
@@ -246,43 +246,6 @@ impl v1alpha1::HiveCluster {
246246
format!("{name}-{role}", name = self.name_any(), role = hive_role)
247247
}
248248

249-
/// Set of functions to define service names on rolegroup level.
250-
/// Headless service for cluster internal purposes only.
251-
// TODO: Move to operator-rs
252-
pub fn rolegroup_headless_service_name(
253-
&self,
254-
rolegroup: &RoleGroupRef<v1alpha1::HiveCluster>,
255-
) -> String {
256-
format!("{name}-headless", name = rolegroup.object_name())
257-
}
258-
259-
/// Headless metrics service exposes Prometheus endpoint only
260-
// TODO: Move to operator-rs
261-
pub fn rolegroup_headless_metrics_service_name(
262-
&self,
263-
rolegroup: &RoleGroupRef<v1alpha1::HiveCluster>,
264-
) -> String {
265-
format!("{name}-metrics", name = rolegroup.object_name())
266-
}
267-
268-
pub fn metrics_ports(&self) -> Vec<ServicePort> {
269-
vec![ServicePort {
270-
name: Some(METRICS_PORT_NAME.to_string()),
271-
port: METRICS_PORT.into(),
272-
protocol: Some("TCP".to_string()),
273-
..ServicePort::default()
274-
}]
275-
}
276-
277-
pub fn service_ports(&self) -> Vec<ServicePort> {
278-
vec![ServicePort {
279-
name: Some(HIVE_PORT_NAME.to_string()),
280-
port: HIVE_PORT.into(),
281-
protocol: Some("TCP".to_string()),
282-
..ServicePort::default()
283-
}]
284-
}
285-
286249
pub fn rolegroup(
287250
&self,
288251
rolegroup_ref: &RoleGroupRef<Self>,

rust/operator-binary/src/main.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ mod kerberos;
77
mod listener;
88
mod operations;
99
mod product_logging;
10+
mod service;
1011

1112
use std::sync::Arc;
1213

Lines changed: 141 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,141 @@
1+
use snafu::{ResultExt, Snafu};
2+
use stackable_operator::{
3+
builder::meta::ObjectMetaBuilder,
4+
commons::product_image_selection::ResolvedProductImage,
5+
k8s_openapi::api::core::v1::{Service, ServicePort, ServiceSpec},
6+
kvp::{Label, Labels},
7+
role_utils::RoleGroupRef,
8+
};
9+
10+
use crate::{
11+
controller::build_recommended_labels,
12+
crd::{APP_NAME, HIVE_PORT, HIVE_PORT_NAME, METRICS_PORT, METRICS_PORT_NAME, v1alpha1},
13+
};
14+
15+
#[derive(Debug, Snafu)]
16+
pub enum Error {
17+
#[snafu(display("object is missing metadata to build owner reference"))]
18+
ObjectMissingMetadataForOwnerRef {
19+
source: stackable_operator::builder::meta::Error,
20+
},
21+
#[snafu(display("failed to build Metadata"))]
22+
MetadataBuild {
23+
source: stackable_operator::builder::meta::Error,
24+
},
25+
#[snafu(display("failed to build Labels"))]
26+
LabelBuild {
27+
source: stackable_operator::kvp::LabelError,
28+
},
29+
}
30+
31+
/// The rolegroup [`Service`] is a headless service that allows direct access to the instances of a certain rolegroup
32+
///
33+
/// This is mostly useful for internal communication between peers, or for clients that perform client-side load balancing.
34+
pub fn build_rolegroup_headless_service(
35+
hive: &v1alpha1::HiveCluster,
36+
resolved_product_image: &ResolvedProductImage,
37+
rolegroup: &RoleGroupRef<v1alpha1::HiveCluster>,
38+
) -> Result<Service, Error> {
39+
let headless_service = Service {
40+
metadata: ObjectMetaBuilder::new()
41+
.name_and_namespace(hive)
42+
// TODO: Use method on RoleGroupRef once op-rs is released
43+
.name(rolegroup_headless_service_name(rolegroup))
44+
.ownerreference_from_resource(hive, None, Some(true))
45+
.context(ObjectMissingMetadataForOwnerRefSnafu)?
46+
.with_recommended_labels(build_recommended_labels(
47+
hive,
48+
&resolved_product_image.app_version_label,
49+
&rolegroup.role,
50+
&rolegroup.role_group,
51+
))
52+
.context(MetadataBuildSnafu)?
53+
.build(),
54+
spec: Some(ServiceSpec {
55+
// Internal communication does not need to be exposed
56+
type_: Some("ClusterIP".to_string()),
57+
cluster_ip: Some("None".to_string()),
58+
// Expecting same ports as on listener service, just as a headless, internal service
59+
ports: Some(service_ports()),
60+
selector: Some(
61+
Labels::role_group_selector(hive, APP_NAME, &rolegroup.role, &rolegroup.role_group)
62+
.context(LabelBuildSnafu)?
63+
.into(),
64+
),
65+
publish_not_ready_addresses: Some(true),
66+
..ServiceSpec::default()
67+
}),
68+
status: None,
69+
};
70+
Ok(headless_service)
71+
}
72+
73+
/// The rolegroup metrics [`Service`] is a service that exposes metrics and a prometheus scraping label
74+
pub fn build_rolegroup_metrics_service(
75+
hive: &v1alpha1::HiveCluster,
76+
resolved_product_image: &ResolvedProductImage,
77+
rolegroup: &RoleGroupRef<v1alpha1::HiveCluster>,
78+
) -> Result<Service, Error> {
79+
let metrics_service = Service {
80+
metadata: ObjectMetaBuilder::new()
81+
.name_and_namespace(hive)
82+
// TODO: Use method on RoleGroupRef once op-rs is released
83+
.name(rolegroup_metrics_service_name(rolegroup))
84+
.ownerreference_from_resource(hive, None, Some(true))
85+
.context(ObjectMissingMetadataForOwnerRefSnafu)?
86+
.with_recommended_labels(build_recommended_labels(
87+
hive,
88+
&resolved_product_image.app_version_label,
89+
&rolegroup.role,
90+
&rolegroup.role_group,
91+
))
92+
.context(MetadataBuildSnafu)?
93+
.with_label(Label::try_from(("prometheus.io/scrape", "true")).context(LabelBuildSnafu)?)
94+
.build(),
95+
spec: Some(ServiceSpec {
96+
// Internal communication does not need to be exposed
97+
type_: Some("ClusterIP".to_string()),
98+
cluster_ip: Some("None".to_string()),
99+
ports: Some(metrics_ports()),
100+
selector: Some(
101+
Labels::role_group_selector(hive, APP_NAME, &rolegroup.role, &rolegroup.role_group)
102+
.context(LabelBuildSnafu)?
103+
.into(),
104+
),
105+
publish_not_ready_addresses: Some(true),
106+
..ServiceSpec::default()
107+
}),
108+
status: None,
109+
};
110+
Ok(metrics_service)
111+
}
112+
113+
/// Headless service for cluster internal purposes only.
114+
// TODO: Move to operator-rs
115+
pub fn rolegroup_headless_service_name(rolegroup: &RoleGroupRef<v1alpha1::HiveCluster>) -> String {
116+
format!("{name}-headless", name = rolegroup.object_name())
117+
}
118+
119+
/// Headless metrics service exposes Prometheus endpoint only
120+
// TODO: Move to operator-rs
121+
pub fn rolegroup_metrics_service_name(rolegroup: &RoleGroupRef<v1alpha1::HiveCluster>) -> String {
122+
format!("{name}-metrics", name = rolegroup.object_name())
123+
}
124+
125+
fn metrics_ports() -> Vec<ServicePort> {
126+
vec![ServicePort {
127+
name: Some(METRICS_PORT_NAME.to_string()),
128+
port: METRICS_PORT.into(),
129+
protocol: Some("TCP".to_string()),
130+
..ServicePort::default()
131+
}]
132+
}
133+
134+
fn service_ports() -> Vec<ServicePort> {
135+
vec![ServicePort {
136+
name: Some(HIVE_PORT_NAME.to_string()),
137+
port: HIVE_PORT.into(),
138+
protocol: Some("TCP".to_string()),
139+
..ServicePort::default()
140+
}]
141+
}

0 commit comments

Comments
 (0)