Skip to content

Commit b83120d

Browse files
committed
add missing file
1 parent 8658f8b commit b83120d

1 file changed

Lines changed: 124 additions & 0 deletions

File tree

Lines changed: 124 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,124 @@
1+
use std::collections::BTreeMap;
2+
3+
use snafu::{ResultExt, Snafu};
4+
use stackable_operator::{
5+
builder::meta::ObjectMetaBuilder,
6+
k8s_openapi::api::core::v1::{Service, ServicePort, ServiceSpec},
7+
kvp::{Label, ObjectLabels},
8+
role_utils::RoleGroupRef,
9+
};
10+
11+
use crate::crd::{HTTPS_PORT, HTTPS_PORT_NAME, METRICS_PORT, METRICS_PORT_NAME, v1alpha1};
12+
13+
const METRICS_SERVICE_SUFFIX: &str = "metrics";
14+
const HEADLESS_SERVICE_SUFFIX: &str = "headless";
15+
16+
#[derive(Snafu, Debug)]
17+
pub enum Error {
18+
#[snafu(display("object is missing metadata to build owner reference"))]
19+
ObjectMissingMetadataForOwnerRef {
20+
source: stackable_operator::builder::meta::Error,
21+
},
22+
23+
#[snafu(display("failed to build Metadata"))]
24+
MetadataBuild {
25+
source: stackable_operator::builder::meta::Error,
26+
},
27+
28+
#[snafu(display("failed to build Labels"))]
29+
LabelBuild {
30+
source: stackable_operator::kvp::LabelError,
31+
},
32+
}
33+
34+
/// The rolegroup headless [`Service`] is a service that allows direct access to the instances of a certain rolegroup
35+
/// This is mostly useful for internal communication between peers, or for clients that perform client-side load balancing.
36+
pub fn build_rolegroup_headless_service(
37+
nifi: &v1alpha1::NifiCluster,
38+
role_group_ref: &RoleGroupRef<v1alpha1::NifiCluster>,
39+
object_labels: ObjectLabels<v1alpha1::NifiCluster>,
40+
selector: BTreeMap<String, String>,
41+
) -> Result<Service, Error> {
42+
Ok(Service {
43+
metadata: ObjectMetaBuilder::new()
44+
.name_and_namespace(nifi)
45+
.name(rolegroup_headless_service_name(
46+
&role_group_ref.object_name(),
47+
))
48+
.ownerreference_from_resource(nifi, None, Some(true))
49+
.context(ObjectMissingMetadataForOwnerRefSnafu)?
50+
.with_recommended_labels(object_labels)
51+
.context(MetadataBuildSnafu)?
52+
.build(),
53+
spec: Some(ServiceSpec {
54+
// Internal communication does not need to be exposed
55+
type_: Some("ClusterIP".to_string()),
56+
cluster_ip: Some("None".to_string()),
57+
ports: Some(headless_service_ports()),
58+
selector: Some(selector),
59+
publish_not_ready_addresses: Some(true),
60+
..ServiceSpec::default()
61+
}),
62+
status: None,
63+
})
64+
}
65+
66+
/// The rolegroup metrics [`Service`] is a service that exposes metrics and a prometheus scraping label.
67+
pub fn build_rolegroup_metrics_service(
68+
nifi: &v1alpha1::NifiCluster,
69+
role_group_ref: &RoleGroupRef<v1alpha1::NifiCluster>,
70+
object_labels: ObjectLabels<v1alpha1::NifiCluster>,
71+
selector: BTreeMap<String, String>,
72+
) -> Result<Service, Error> {
73+
Ok(Service {
74+
metadata: ObjectMetaBuilder::new()
75+
.name_and_namespace(nifi)
76+
.name(rolegroup_metrics_service_name(
77+
&role_group_ref.object_name(),
78+
))
79+
.ownerreference_from_resource(nifi, None, Some(true))
80+
.context(ObjectMissingMetadataForOwnerRefSnafu)?
81+
.with_recommended_labels(object_labels)
82+
.context(MetadataBuildSnafu)?
83+
.with_label(Label::try_from(("prometheus.io/scrape", "true")).context(LabelBuildSnafu)?)
84+
.build(),
85+
spec: Some(ServiceSpec {
86+
// Internal communication does not need to be exposed
87+
type_: Some("ClusterIP".to_string()),
88+
cluster_ip: Some("None".to_string()),
89+
ports: Some(metrics_service_ports()),
90+
selector: Some(selector),
91+
publish_not_ready_addresses: Some(true),
92+
..ServiceSpec::default()
93+
}),
94+
status: None,
95+
})
96+
}
97+
98+
fn headless_service_ports() -> Vec<ServicePort> {
99+
vec![ServicePort {
100+
name: Some(HTTPS_PORT_NAME.into()),
101+
port: HTTPS_PORT.into(),
102+
protocol: Some("TCP".to_string()),
103+
..ServicePort::default()
104+
}]
105+
}
106+
107+
fn metrics_service_ports() -> Vec<ServicePort> {
108+
vec![ServicePort {
109+
name: Some(METRICS_PORT_NAME.to_string()),
110+
port: METRICS_PORT.into(),
111+
protocol: Some("TCP".to_string()),
112+
..ServicePort::default()
113+
}]
114+
}
115+
116+
/// Returns the metrics rolegroup service name `<cluster>-<role>-<rolegroup>-<METRICS_SERVICE_SUFFIX>`.
117+
fn rolegroup_metrics_service_name(role_group_ref_object_name: &str) -> String {
118+
format!("{role_group_ref_object_name}-{METRICS_SERVICE_SUFFIX}")
119+
}
120+
121+
/// Returns the headless rolegroup service name `<cluster>-<role>-<rolegroup>-<HEADLESS_SERVICE_SUFFIX>`.
122+
pub fn rolegroup_headless_service_name(role_group_ref_object_name: &str) -> String {
123+
format!("{role_group_ref_object_name}-{HEADLESS_SERVICE_SUFFIX}")
124+
}

0 commit comments

Comments
 (0)