44
55use std:: str:: FromStr ;
66
7+ use snafu:: { OptionExt , ResultExt , Snafu } ;
78use 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+
924pub mod git_sync;
1025pub mod graceful_shutdown;
1126pub 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.
2843pub const NIFI_CONFIG_DIRECTORY : & str = "/stackable/nifi/conf" ;
2944pub 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+ }
0 commit comments