@@ -15,16 +15,18 @@ use csi_server::{
1515use futures:: { FutureExt , TryFutureExt , TryStreamExt } ;
1616use stackable_operator:: {
1717 self , YamlSchema ,
18- cli:: { CommonOptions , MaintenanceOptions , OperatorEnvironmentOptions } ,
18+ cli:: { Command , CommonOptions , MaintenanceOptions , OperatorEnvironmentOptions } ,
19+ client:: Client ,
1920 crd:: listener:: {
2021 Listener , ListenerClass , ListenerClassVersion , ListenerVersion , PodListeners ,
21- PodListenersVersion ,
22+ PodListenersVersion , v1alpha1 ,
2223 } ,
2324 eos:: EndOfSupportChecker ,
2425 shared:: yaml:: SerializeOptions ,
2526 telemetry:: Tracing ,
2627 utils:: signal:: SignalWatcher ,
2728} ;
29+ use tokio:: sync:: oneshot;
2830use tokio_stream:: wrappers:: UnixListenerStream ;
2931use tonic:: transport:: Server ;
3032use utils:: unix_stream:: { TonicUnixStream , uds_bind_private} ;
@@ -42,14 +44,14 @@ const FIELD_MANAGER: &str = "listener-operator";
4244
4345#[ derive( clap:: Parser ) ]
4446#[ clap( author, version) ]
45- struct Opts {
47+ struct Cli {
4648 #[ clap( subcommand) ]
47- cmd : stackable_operator :: cli :: Command < ListenerOperatorRun > ,
49+ cmd : Command < ListenerOperatorRun > ,
4850}
4951
5052#[ derive( clap:: Parser ) ]
5153struct ListenerOperatorRun {
52- #[ clap ( long, env) ]
54+ #[ arg ( long, env) ]
5355 csi_endpoint : PathBuf ,
5456
5557 #[ clap( subcommand) ]
@@ -70,12 +72,34 @@ struct ListenerOperatorRun {
7072#[ derive( Debug , clap:: Parser , strum:: AsRefStr , strum:: Display ) ]
7173enum RunMode {
7274 /// CSI Controller Service
73- Controller ,
75+ Controller ( ControllerArguments ) ,
7476
7577 /// CSI Node Service
7678 Node ,
7779}
7880
81+ #[ derive( Debug , clap:: Args ) ]
82+ struct ControllerArguments {
83+ #[ arg( long, env, default_value_t) ]
84+ listener_class_preset : ListenerClassPreset ,
85+ }
86+
87+ #[ derive( Clone , Debug , Default , clap:: Parser , strum:: Display , strum:: EnumString ) ]
88+ #[ strum( serialize_all = "kebab-case" ) ]
89+ enum ListenerClassPreset {
90+ /// Deploys no listener class preset.
91+ None ,
92+
93+ /// Deploys listener classes for environments in which pods can move freely
94+ /// between nodes. This is common for many managed cloud environments.
95+ #[ default]
96+ EphemeralNodes ,
97+
98+ /// Deploys listener classes for environments with reliable, long-living
99+ /// nodes and pods don't move between nodes.
100+ StableNodes ,
101+ }
102+
79103mod built_info {
80104 include ! ( concat!( env!( "OUT_DIR" ) , "/built.rs" ) ) ;
81105}
@@ -85,17 +109,17 @@ pub const ENV_VAR_CONSOLE_LOG: &str = "LISTENER_OPERATOR_LOG";
85109
86110#[ tokio:: main]
87111async fn main ( ) -> anyhow:: Result < ( ) > {
88- let opts = Opts :: parse ( ) ;
112+ let opts = Cli :: parse ( ) ;
89113 match opts. cmd {
90- stackable_operator :: cli :: Command :: Crd => {
114+ Command :: Crd => {
91115 ListenerClass :: merged_crd ( ListenerClassVersion :: V1Alpha1 ) ?
92116 . print_yaml_schema ( built_info:: PKG_VERSION , SerializeOptions :: default ( ) ) ?;
93117 Listener :: merged_crd ( ListenerVersion :: V1Alpha1 ) ?
94118 . print_yaml_schema ( built_info:: PKG_VERSION , SerializeOptions :: default ( ) ) ?;
95119 PodListeners :: merged_crd ( PodListenersVersion :: V1Alpha1 ) ?
96120 . print_yaml_schema ( built_info:: PKG_VERSION , SerializeOptions :: default ( ) ) ?;
97121 }
98- stackable_operator :: cli :: Command :: Run ( ListenerOperatorRun {
122+ Command :: Run ( ListenerOperatorRun {
99123 operator_environment,
100124 csi_endpoint,
101125 maintenance,
@@ -154,29 +178,47 @@ async fn main() -> anyhow::Result<()> {
154178 )
155179 . add_service ( IdentityServer :: new ( ListenerOperatorIdentity ) ) ;
156180
157- let webhook_server = create_webhook_server (
158- & operator_environment,
159- maintenance. disable_crd_maintenance ,
160- client. as_kube_client ( ) ,
161- )
162- . await ?;
181+ match mode {
182+ RunMode :: Controller ( ControllerArguments {
183+ listener_class_preset,
184+ } ) => {
185+ let ( webhook_server, initial_reconcile_rx) = create_webhook_server (
186+ & operator_environment,
187+ maintenance. disable_crd_maintenance ,
188+ client. as_kube_client ( ) ,
189+ )
190+ . await ?;
163191
164- let webhook_server = webhook_server
165- . run ( sigterm_watcher. handle ( ) )
166- . map_err ( |err| anyhow ! ( err) . context ( "failed to run webhook server" ) ) ;
192+ let webhook_server = webhook_server
193+ . run ( sigterm_watcher. handle ( ) )
194+ . map_err ( |err| anyhow ! ( err) . context ( "failed to run webhook server" ) ) ;
195+
196+ let listener_classes = create_listener_classes (
197+ initial_reconcile_rx,
198+ listener_class_preset,
199+ client. clone ( ) ,
200+ )
201+ . map_err ( |err| {
202+ anyhow ! ( err) . context ( "failed to apply listener classes selected by preset" )
203+ } ) ;
167204
168- match mode {
169- RunMode :: Controller => {
170205 let csi_server = csi_server
171206 . add_service ( ControllerServer :: new ( ListenerOperatorController {
172207 client : client. clone ( ) ,
173208 } ) )
174209 . serve_with_incoming_shutdown ( csi_listener, sigterm_watcher. handle ( ) )
175210 . map_err ( |err| anyhow ! ( err) . context ( "failed to run csi server" ) ) ;
211+
176212 let controller =
177213 listener_controller:: run ( client, sigterm_watcher. handle ( ) ) . map ( anyhow:: Ok ) ;
178214
179- futures:: try_join!( csi_server, controller, eos_checker, webhook_server) ?;
215+ futures:: try_join!(
216+ listener_classes,
217+ webhook_server,
218+ eos_checker,
219+ csi_server,
220+ controller,
221+ ) ?;
180222 }
181223 RunMode :: Node => {
182224 let node_name = & common. cluster_info . kubernetes_node_name ;
@@ -188,10 +230,38 @@ async fn main() -> anyhow::Result<()> {
188230 . serve_with_incoming_shutdown ( csi_listener, sigterm_watcher. handle ( ) )
189231 . map_err ( |err| anyhow ! ( err) . context ( "failed to run csi server" ) ) ;
190232
191- futures:: try_join!( csi_server, eos_checker, webhook_server ) ?;
233+ futures:: try_join!( csi_server, eos_checker) ?;
192234 }
193235 }
194236 }
195237 }
238+
239+ Ok ( ( ) )
240+ }
241+
242+ async fn create_listener_classes (
243+ initial_reconcile_rx : oneshot:: Receiver < ( ) > ,
244+ listener_class_preset : ListenerClassPreset ,
245+ client : Client ,
246+ ) -> anyhow:: Result < ( ) > {
247+ initial_reconcile_rx. await ?;
248+
249+ tracing:: info!( "applying \" {listener_class_preset}\" listener class preset" ) ;
250+
251+ #[ rustfmt:: skip]
252+ let bytes = match listener_class_preset {
253+ ListenerClassPreset :: None => return Ok ( ( ) ) ,
254+ ListenerClassPreset :: EphemeralNodes => include_bytes ! ( "manifests/ephemeral-nodes.yaml" ) . to_vec ( ) ,
255+ ListenerClassPreset :: StableNodes => include_bytes ! ( "manifests/stable-nodes.yaml" ) . to_vec ( ) ,
256+ } ;
257+
258+ for document in serde_yaml:: Deserializer :: from_slice ( & bytes) {
259+ let class: v1alpha1:: ListenerClass =
260+ serde_yaml:: with:: singleton_map_recursive:: deserialize ( document)
261+ . expect ( "compile-time included listener classes must be valid YAML" ) ;
262+
263+ client. create_if_missing ( & class) . await ?;
264+ }
265+
196266 Ok ( ( ) )
197267}
0 commit comments