Skip to content

Commit 75728c5

Browse files
committed
refactor: Introduce build aggregator
1 parent 5a68359 commit 75728c5

3 files changed

Lines changed: 202 additions & 148 deletions

File tree

rust/operator-binary/src/controller/build.rs

Lines changed: 143 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,8 +4,23 @@
44
55
use std::str::FromStr;
66

7+
use snafu::{OptionExt, ResultExt, Snafu};
78
use stackable_operator::v2::types::{common::Port, operator::RoleGroupName};
89

10+
use crate::{
11+
controller::{
12+
KubernetesResources, ValidatedCluster,
13+
build::resource::{
14+
config_map::build_rolegroup_config_map,
15+
listener::{build_group_listener, group_listener_name},
16+
pdb::build_pdb,
17+
service::{build_rolegroup_headless_service, build_rolegroup_metrics_service},
18+
statefulset::build_node_rolegroup_statefulset,
19+
},
20+
},
21+
crd::NifiRole,
22+
};
23+
924
pub mod git_sync;
1025
pub mod graceful_shutdown;
1126
pub mod jvm;
@@ -27,3 +42,131 @@ pub const BALANCE_PORT: Port = Port(6243);
2742
// Filesystem paths shared by multiple builders. Single-consumer paths live in their builder.
2843
pub const NIFI_CONFIG_DIRECTORY: &str = "/stackable/nifi/conf";
2944
pub const NIFI_PYTHON_WORKING_DIRECTORY: &str = "/nifi-python-working-directory";
45+
46+
#[derive(Snafu, Debug)]
47+
pub enum Error {
48+
#[snafu(display("NifiCluster has no nodes role defined"))]
49+
NoNodesDefined,
50+
51+
#[snafu(display("failed to build ConfigMap for role group {role_group}"))]
52+
ConfigMap {
53+
source: resource::config_map::Error,
54+
role_group: RoleGroupName,
55+
},
56+
57+
#[snafu(display("failed to build StatefulSet for role group {role_group}"))]
58+
StatefulSet {
59+
source: resource::statefulset::Error,
60+
role_group: RoleGroupName,
61+
},
62+
}
63+
64+
/// Builds every Kubernetes resource for the given validated cluster.
65+
///
66+
/// Does not need a Kubernetes client: every reference to another Kubernetes resource is already
67+
/// dereferenced and validated by this point, so the errors returned here are resource-assembly
68+
/// failures only.
69+
///
70+
/// `service_account_name` is the name of the RBAC `ServiceAccount` the role-group Pods run under
71+
/// (RBAC resources are built and applied separately, in the reconcile step).
72+
pub fn build(
73+
cluster: &ValidatedCluster,
74+
service_account_name: &str,
75+
) -> Result<KubernetesResources, Error> {
76+
let mut stateful_sets = vec![];
77+
let mut services = vec![];
78+
let mut listeners = vec![];
79+
let mut config_maps = vec![];
80+
let mut pod_disruption_budgets = vec![];
81+
82+
// NiFi has a single role (`node`), which must always be present.
83+
let nifi_role = NifiRole::Node;
84+
let node_role_group_configs = cluster
85+
.role_group_configs
86+
.get(&nifi_role)
87+
.context(NoNodesDefinedSnafu)?;
88+
89+
// Role-level resources (one per role): the PodDisruptionBudget and the group Listener.
90+
let role_config = &cluster.role_config;
91+
if let Some(pdb) = build_pdb(&role_config.pdb, cluster, &nifi_role) {
92+
pod_disruption_budgets.push(pdb);
93+
}
94+
listeners.push(build_group_listener(
95+
cluster,
96+
role_config.listener_class.clone(),
97+
group_listener_name(cluster, &nifi_role.to_string()),
98+
));
99+
100+
for (role_group_name, rg) in node_role_group_configs {
101+
services.push(build_rolegroup_headless_service(cluster, role_group_name));
102+
services.push(build_rolegroup_metrics_service(cluster, role_group_name));
103+
104+
config_maps.push(
105+
build_rolegroup_config_map(cluster, role_group_name, rg).context(ConfigMapSnafu {
106+
role_group: role_group_name.clone(),
107+
})?,
108+
);
109+
110+
let effective_replicas = rg.replicas.map(i32::from);
111+
stateful_sets.push(
112+
build_node_rolegroup_statefulset(
113+
cluster,
114+
role_group_name,
115+
rg,
116+
effective_replicas,
117+
service_account_name,
118+
)
119+
.context(StatefulSetSnafu {
120+
role_group: role_group_name.clone(),
121+
})?,
122+
);
123+
}
124+
125+
Ok(KubernetesResources {
126+
stateful_sets,
127+
services,
128+
listeners,
129+
config_maps,
130+
pod_disruption_budgets,
131+
})
132+
}
133+
134+
#[cfg(test)]
135+
mod tests {
136+
use stackable_operator::kube::Resource;
137+
138+
use super::{build, properties::test_support::minimal_validated_cluster};
139+
140+
fn sorted_names(resources: &[impl Resource]) -> Vec<&str> {
141+
let mut names: Vec<&str> = resources
142+
.iter()
143+
.filter_map(|resource| resource.meta().name.as_deref())
144+
.collect();
145+
names.sort();
146+
names
147+
}
148+
149+
#[test]
150+
fn build_produces_expected_resources() {
151+
let cluster = minimal_validated_cluster();
152+
let resources = build(&cluster, "simple-nifi-serviceaccount").expect("build succeeds");
153+
154+
// The minimal fixture has a single `default` role group for the `node` role.
155+
assert_eq!(
156+
sorted_names(&resources.stateful_sets),
157+
["simple-nifi-node-default"]
158+
);
159+
assert_eq!(
160+
sorted_names(&resources.config_maps),
161+
["simple-nifi-node-default"]
162+
);
163+
// One headless and one metrics Service per role group.
164+
assert_eq!(resources.services.len(), 2);
165+
// One group Listener and one PDB for the single `node` role.
166+
assert_eq!(sorted_names(&resources.listeners), ["simple-nifi-node"]);
167+
assert_eq!(
168+
sorted_names(&resources.pod_disruption_budgets),
169+
["simple-nifi-node"]
170+
);
171+
}
172+
}

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

Lines changed: 18 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -11,8 +11,15 @@ use stackable_operator::{
1111
product_image_selection::ResolvedProductImage,
1212
resources::{NoRuntimeLimits, Resources},
1313
},
14-
crd::git_sync,
15-
k8s_openapi::{api::core::v1::Volume, apimachinery::pkg::apis::meta::v1::ObjectMeta},
14+
crd::{git_sync, listener},
15+
k8s_openapi::{
16+
api::{
17+
apps::v1::StatefulSet,
18+
core::v1::{ConfigMap, Service, Volume},
19+
policy::v1::PodDisruptionBudget,
20+
},
21+
apimachinery::pkg::apis::meta::v1::ObjectMeta,
22+
},
1623
kube::Resource,
1724
kvp::Labels,
1825
shared::time::Duration,
@@ -48,6 +55,15 @@ pub(crate) mod build;
4855
pub(crate) mod dereference;
4956
pub(crate) mod validate;
5057

58+
/// Every Kubernetes resource produced by the [`build`] step.
59+
pub struct KubernetesResources {
60+
pub stateful_sets: Vec<StatefulSet>,
61+
pub services: Vec<Service>,
62+
pub listeners: Vec<listener::v1alpha1::Listener>,
63+
pub config_maps: Vec<ConfigMap>,
64+
pub pod_disruption_budgets: Vec<PodDisruptionBudget>,
65+
}
66+
5167
/// A validated, merged (default <- role <- role-group) NiFi rolegroup config.
5268
pub type NifiRoleGroupConfig =
5369
RoleGroupConfig<ValidatedNifiConfig, JavaCommonConfig, v1alpha1::NifiConfigOverrides>;

0 commit comments

Comments
 (0)