Skip to main content

mz_orchestratord/controller/
console.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10use k8s_controller::{Outcome, TraceMetadata};
11use k8s_openapi::{
12    api::{
13        apps::v1::{Deployment, DeploymentSpec},
14        core::v1::{
15            Affinity, Capabilities, ConfigMap, ConfigMapVolumeSource, Container, ContainerPort,
16            EnvVar, HTTPGetAction, KeyToPath, PodSecurityContext, PodSpec, PodTemplateSpec, Probe,
17            ResourceRequirements, SeccompProfile, SecretVolumeSource, SecurityContext, Service,
18            ServicePort, ServiceSpec, Toleration, Volume, VolumeMount,
19        },
20        networking::v1::{
21            IPBlock, NetworkPolicy, NetworkPolicyIngressRule, NetworkPolicyPeer, NetworkPolicyPort,
22            NetworkPolicySpec,
23        },
24    },
25    apimachinery::pkg::{
26        apis::meta::v1::{Condition, LabelSelector, Time},
27        util::intstr::IntOrString,
28    },
29    jiff::Timestamp,
30};
31use kube::{
32    Api, Client, Resource, ResourceExt,
33    api::{DeleteParams, ObjectMeta, PostParams},
34    runtime::{conditions::is_deployment_completed, controller::Action, wait::await_condition},
35};
36use maplit::btreemap;
37use serde::Serialize;
38use tracing::{trace, warn};
39
40use crate::{
41    Error,
42    k8s::{apply_resource, get_resource, recommended_k8s_labels, replace_resource},
43    tls::{DefaultCertificateSpecs, create_certificate, issuer_ref_defined},
44};
45use mz_cloud_resources::crd::{
46    ManagedResource,
47    console::v1alpha1::{Console, HttpConnectionScheme},
48    generated::cert_manager::certificates::{Certificate, CertificatePrivateKeyAlgorithm},
49};
50use mz_orchestrator_kubernetes::KubernetesImagePullPolicy;
51use mz_ore::{cli::KeyValueArg, instrument};
52use mz_server_core::listeners::AuthenticatorKind;
53
54/// The name identifying this controller in its reconciliation metrics and in
55/// the reporter of the events it publishes.
56pub const CONTROLLER_NAME: &str = "console";
57
58#[derive(Clone)]
59pub struct Config {
60    pub enable_security_context: bool,
61    pub enable_prometheus_scrape_annotations: bool,
62
63    pub image_pull_policy: KubernetesImagePullPolicy,
64    pub scheduler_name: Option<String>,
65    pub console_node_selector: Vec<KeyValueArg<String, String>>,
66    pub console_affinity: Option<Affinity>,
67    pub console_tolerations: Option<Vec<Toleration>>,
68    pub console_default_resources: Option<ResourceRequirements>,
69    pub network_policies_ingress_enabled: bool,
70    pub network_policies_ingress_cidrs: Vec<String>,
71
72    pub default_certificate_specs: DefaultCertificateSpecs,
73
74    pub console_http_port: u16,
75    pub balancerd_http_port: u16,
76}
77
78#[derive(Serialize)]
79struct AppConfig {
80    version: String,
81    auth: AppConfigAuth,
82    #[serde(default, skip_serializing_if = "Option::is_none")]
83    balancerd_dns_names: Option<Vec<String>>,
84}
85
86#[derive(Serialize)]
87struct AppConfigAuth {
88    mode: AuthenticatorKind,
89}
90
91pub struct Context {
92    config: Config,
93}
94
95impl Context {
96    pub fn new(config: Config) -> Self {
97        Self { config }
98    }
99
100    async fn sync_deployment_status(
101        &self,
102        client: &Client,
103        console: &Console,
104    ) -> Result<(), Error> {
105        let namespace = console.namespace();
106        let console_api: Api<Console> = Api::namespaced(client.clone(), &namespace);
107        let deployment_api: Api<Deployment> = Api::namespaced(client.clone(), &namespace);
108
109        let Some(deployment) = get_resource(&deployment_api, &console.deployment_name()).await?
110        else {
111            return Ok(());
112        };
113
114        let Some(deployment_conditions) = &deployment
115            .status
116            .as_ref()
117            .and_then(|status| status.conditions.as_ref())
118        else {
119            // if the deployment doesn't have any conditions set yet, there
120            // is nothing to sync
121            return Ok(());
122        };
123
124        let ready = deployment_conditions
125            .iter()
126            .any(|condition| condition.type_ == "Available" && condition.status == "True");
127        let ready_str = if ready { "True" } else { "False" };
128
129        let mut status = console.status.clone().unwrap();
130        if status
131            .conditions
132            .iter()
133            .any(|condition| condition.type_ == "Ready" && condition.status == ready_str)
134        {
135            // if the deployment status is already set correctly, we don't
136            // need to set it again (this prevents us from getting stuck in
137            // a reconcile loop)
138            return Ok(());
139        }
140
141        status.conditions = vec![Condition {
142            type_: "Ready".to_string(),
143            status: ready_str.to_string(),
144            last_transition_time: Time(Timestamp::now()),
145            message: format!(
146                "console deployment is{} ready",
147                if ready { "" } else { " not" }
148            ),
149            observed_generation: None,
150            reason: "DeploymentStatus".to_string(),
151        }];
152        let mut new_console = console.clone();
153        new_console.status = Some(status);
154
155        console_api
156            .replace_status(
157                &console.name_unchecked(),
158                &PostParams::default(),
159                &new_console,
160            )
161            .await?;
162
163        Ok(())
164    }
165
166    fn create_network_policies(&self, console: &Console) -> Vec<NetworkPolicy> {
167        let mut network_policies = Vec::new();
168        if self.config.network_policies_ingress_enabled {
169            let console_label_selector = LabelSelector {
170                match_labels: Some(
171                    console
172                        .default_labels()
173                        .into_iter()
174                        .chain([("materialize.cloud/app".to_owned(), console.app_name())])
175                        .collect(),
176                ),
177                ..Default::default()
178            };
179            network_policies.extend([NetworkPolicy {
180                metadata: console.managed_resource_meta(console.name_prefixed("console-ingress")),
181                spec: Some(NetworkPolicySpec {
182                    ingress: Some(vec![NetworkPolicyIngressRule {
183                        from: Some(
184                            self.config
185                                .network_policies_ingress_cidrs
186                                .iter()
187                                .map(|cidr| NetworkPolicyPeer {
188                                    ip_block: Some(IPBlock {
189                                        cidr: cidr.to_owned(),
190                                        except: None,
191                                    }),
192                                    ..Default::default()
193                                })
194                                .collect(),
195                        ),
196                        ports: Some(vec![NetworkPolicyPort {
197                            port: Some(IntOrString::Int(self.config.console_http_port.into())),
198                            protocol: Some("TCP".to_string()),
199                            ..Default::default()
200                        }]),
201                        ..Default::default()
202                    }]),
203                    pod_selector: Some(console_label_selector),
204                    policy_types: Some(vec!["Ingress".to_owned()]),
205                    ..Default::default()
206                }),
207            }]);
208        }
209        network_policies
210    }
211
212    fn create_console_external_certificate(
213        &self,
214        console: &Console,
215    ) -> anyhow::Result<Option<Certificate>> {
216        create_certificate(
217            self.config
218                .default_certificate_specs
219                .console_external
220                .clone(),
221            console,
222            console.spec.external_certificate_spec.clone(),
223            console.external_certificate_name(),
224            console.external_certificate_secret_name(),
225            None,
226            CertificatePrivateKeyAlgorithm::Ecdsa,
227            Some(256),
228        )
229    }
230
231    fn create_console_app_configmap_object(&self, console: &Console) -> ConfigMap {
232        let balancerd_dns_names = console.spec.balancerd.dns_names.clone();
233        let version: String = console
234            .spec
235            .console_image_ref
236            .rsplitn(2, ':')
237            .next()
238            .expect("at least one chunk, even if empty")
239            .to_owned();
240        let app_config_json = serde_json::to_string(&AppConfig {
241            version,
242            balancerd_dns_names,
243            auth: AppConfigAuth {
244                mode: console.spec.authenticator_kind,
245            },
246        })
247        .expect("known valid");
248        ConfigMap {
249            binary_data: None,
250            data: Some(btreemap! {
251                "app-config.json".to_owned() => app_config_json,
252            }),
253            immutable: None,
254            metadata: console.managed_resource_meta(console.configmap_name()),
255        }
256    }
257
258    fn create_console_deployment_object(&self, console: &Console) -> Deployment {
259        let mut pod_template_labels = console.default_labels();
260        pod_template_labels.insert(
261            "materialize.cloud/name".to_owned(),
262            console.deployment_name(),
263        );
264        pod_template_labels.insert("app".to_owned(), "console".to_string());
265        pod_template_labels.insert("materialize.cloud/app".to_owned(), console.app_name());
266
267        let ports = vec![ContainerPort {
268            container_port: self.config.console_http_port.into(),
269            name: Some("http".into()),
270            protocol: Some("TCP".into()),
271            ..Default::default()
272        }];
273
274        let scheme = match console.spec.balancerd.scheme {
275            HttpConnectionScheme::Http => "http",
276            HttpConnectionScheme::Https => "https",
277        };
278        let mut env = vec![EnvVar {
279            name: "MZ_ENDPOINT".to_string(),
280            value: Some(format!(
281                "{}://{}.{}.svc.cluster.local:{}",
282                scheme,
283                console.spec.balancerd.service_name,
284                console.spec.balancerd.namespace,
285                self.config.balancerd_http_port,
286            )),
287            ..Default::default()
288        }];
289        let mut volumes = vec![Volume {
290            name: "app-config".to_string(),
291            config_map: Some(ConfigMapVolumeSource {
292                name: console.configmap_name(),
293                default_mode: Some(256),
294                optional: Some(false),
295                items: Some(vec![KeyToPath {
296                    key: "app-config.json".to_string(),
297                    path: "app-config.json".to_string(),
298                    ..Default::default()
299                }]),
300            }),
301            ..Default::default()
302        }];
303        let mut volume_mounts = vec![VolumeMount {
304            name: "app-config".to_string(),
305            mount_path: "/usr/share/nginx/html/app-config".to_string(),
306            ..Default::default()
307        }];
308
309        let scheme = if issuer_ref_defined(
310            &self.config.default_certificate_specs.console_external,
311            &console.spec.external_certificate_spec,
312        ) {
313            volumes.push(Volume {
314                name: "external-certificate".to_owned(),
315                secret: Some(SecretVolumeSource {
316                    default_mode: Some(0o400),
317                    secret_name: Some(console.external_certificate_secret_name()),
318                    items: None,
319                    optional: Some(false),
320                }),
321                ..Default::default()
322            });
323            volume_mounts.push(VolumeMount {
324                name: "external-certificate".to_owned(),
325                mount_path: "/nginx/tls".to_owned(),
326                read_only: Some(true),
327                ..Default::default()
328            });
329            env.push(EnvVar {
330                name: "MZ_NGINX_LISTENER_CONFIG".to_string(),
331                value: Some(format!(
332                    "listen {} ssl;
333ssl_certificate /nginx/tls/tls.crt;
334ssl_certificate_key /nginx/tls/tls.key;",
335                    self.config.console_http_port
336                )),
337                ..Default::default()
338            });
339            Some("HTTPS".to_owned())
340        } else {
341            env.push(EnvVar {
342                name: "MZ_NGINX_LISTENER_CONFIG".to_string(),
343                value: Some(format!("listen {};", self.config.console_http_port)),
344                ..Default::default()
345            });
346            Some("HTTP".to_owned())
347        };
348
349        let probe = Probe {
350            http_get: Some(HTTPGetAction {
351                path: Some("/".to_string()),
352                port: IntOrString::Int(self.config.console_http_port.into()),
353                scheme,
354                ..Default::default()
355            }),
356            ..Default::default()
357        };
358
359        let security_context = if self.config.enable_security_context {
360            // Since we want to adhere to the most restrictive security context, all
361            // of these fields have to be set how they are.
362            // See https://kubernetes.io/docs/concepts/security/pod-security-standards/#restricted
363            Some(SecurityContext {
364                run_as_non_root: Some(true),
365                capabilities: Some(Capabilities {
366                    drop: Some(vec!["ALL".to_string()]),
367                    ..Default::default()
368                }),
369                seccomp_profile: Some(SeccompProfile {
370                    type_: "RuntimeDefault".to_string(),
371                    ..Default::default()
372                }),
373                allow_privilege_escalation: Some(false),
374                ..Default::default()
375            })
376        } else {
377            None
378        };
379
380        let container = Container {
381            name: "console".to_owned(),
382            image: Some(console.spec.console_image_ref.clone()),
383            image_pull_policy: Some(self.config.image_pull_policy.to_string()),
384            ports: Some(ports),
385            env: Some(env),
386            startup_probe: Some(Probe {
387                period_seconds: Some(1),
388                failure_threshold: Some(10),
389                ..probe.clone()
390            }),
391            readiness_probe: Some(Probe {
392                period_seconds: Some(30),
393                failure_threshold: Some(1),
394                ..probe.clone()
395            }),
396            liveness_probe: Some(Probe {
397                period_seconds: Some(30),
398                ..probe.clone()
399            }),
400            resources: console
401                .spec
402                .resource_requirements
403                .clone()
404                .or_else(|| self.config.console_default_resources.clone()),
405            security_context,
406            volume_mounts: Some(volume_mounts),
407            ..Default::default()
408        };
409
410        let match_labels = pod_template_labels.clone();
411        pod_template_labels.extend(recommended_k8s_labels("console".into()));
412
413        let deployment_spec = DeploymentSpec {
414            replicas: Some(console.replicas()),
415            selector: LabelSelector {
416                match_labels: Some(match_labels),
417                ..Default::default()
418            },
419            template: PodTemplateSpec {
420                // not using managed_resource_meta because the pod should be
421                // owned by the deployment, not the materialize instance
422                metadata: Some(ObjectMeta {
423                    labels: Some(pod_template_labels),
424                    ..Default::default()
425                }),
426                spec: Some(PodSpec {
427                    containers: vec![container],
428                    node_selector: Some(
429                        self.config
430                            .console_node_selector
431                            .iter()
432                            .map(|selector| (selector.key.clone(), selector.value.clone()))
433                            .collect(),
434                    ),
435                    affinity: self.config.console_affinity.clone(),
436                    tolerations: self.config.console_tolerations.clone(),
437                    scheduler_name: self.config.scheduler_name.clone(),
438                    volumes: Some(volumes),
439                    security_context: Some(PodSecurityContext {
440                        fs_group: Some(101),
441                        ..Default::default()
442                    }),
443                    ..Default::default()
444                }),
445            },
446            ..Default::default()
447        };
448
449        Deployment {
450            metadata: ObjectMeta {
451                ..console.managed_resource_meta(console.deployment_name())
452            },
453            spec: Some(deployment_spec),
454            status: None,
455        }
456    }
457
458    fn create_console_service_object(&self, console: &Console) -> Service {
459        let selector =
460            btreemap! {"materialize.cloud/name".to_string() => console.deployment_name()};
461
462        let ports = vec![ServicePort {
463            name: Some("http".to_string()),
464            protocol: Some("TCP".to_string()),
465            port: self.config.console_http_port.into(),
466            target_port: Some(IntOrString::Int(self.config.console_http_port.into())),
467            ..Default::default()
468        }];
469
470        let spec = ServiceSpec {
471            type_: Some("ClusterIP".to_string()),
472            cluster_ip: Some("None".to_string()),
473            selector: Some(selector),
474            ports: Some(ports),
475            ..Default::default()
476        };
477
478        Service {
479            metadata: console.managed_resource_meta(console.service_name()),
480            spec: Some(spec),
481            status: None,
482        }
483    }
484
485    // TODO: remove this once everyone is upgraded to an orchestratord
486    // version with the separate console operator
487    async fn fix_deployment(
488        &self,
489        deployment_api: &Api<Deployment>,
490        new_deployment: &Deployment,
491    ) -> Result<(), Error> {
492        let Some(mut existing_deployment) =
493            get_resource(deployment_api, &new_deployment.name_unchecked()).await?
494        else {
495            return Ok(());
496        };
497
498        if existing_deployment.spec.as_ref().unwrap().selector
499            == new_deployment.spec.as_ref().unwrap().selector
500        {
501            return Ok(());
502        }
503
504        warn!("found existing deployment with old label selector, fixing");
505
506        // this is sufficient because the new labels are a superset of the
507        // old labels, so the existing label selector should still be valid
508        existing_deployment
509            .spec
510            .as_mut()
511            .unwrap()
512            .template
513            .metadata
514            .as_mut()
515            .unwrap()
516            .labels = new_deployment
517            .spec
518            .as_ref()
519            .unwrap()
520            .template
521            .metadata
522            .as_ref()
523            .unwrap()
524            .labels
525            .clone();
526
527        // using await_condition is not ideal in a controller loop, but this
528        // is very temporary and will only ever happen once, so this feels
529        // simpler than trying to introduce an entire state machine here
530        replace_resource(deployment_api, &existing_deployment).await?;
531        await_condition(
532            deployment_api.clone(),
533            &existing_deployment.name_unchecked(),
534            |deployment: Option<&Deployment>| {
535                let observed_generation = deployment
536                    .and_then(|deployment| deployment.status.as_ref())
537                    .and_then(|status| status.observed_generation)
538                    .unwrap_or(0);
539                let current_generation = deployment
540                    .and_then(|deployment| deployment.meta().generation)
541                    .unwrap_or(0);
542                let previous_generation = existing_deployment.meta().generation.unwrap_or(0);
543                observed_generation == current_generation
544                    && current_generation > previous_generation
545            },
546        )
547        .await
548        .map_err(|e| anyhow::anyhow!(e))?;
549        await_condition(
550            deployment_api.clone(),
551            &existing_deployment.name_unchecked(),
552            is_deployment_completed(),
553        )
554        .await
555        .map_err(|e| anyhow::anyhow!(e))?;
556
557        // delete the deployment but leave the pods around (via
558        // DeleteParams::orphan)
559        match kube::runtime::wait::delete::delete_and_finalize(
560            deployment_api.clone(),
561            &existing_deployment.name_unchecked(),
562            &DeleteParams::orphan(),
563        )
564        .await
565        {
566            Ok(_) => {}
567            Err(kube::runtime::wait::delete::Error::Delete(kube::Error::Api(e)))
568                if e.code == 404 =>
569            {
570                // the resource already doesn't exist
571            }
572            Err(e) => return Err(anyhow::anyhow!(e).into()),
573        }
574
575        // now, the normal apply of the new deployment (in the main loop)
576        // will take over the existing pods from the old deployment we just
577        // deleted, since we already updated the pod labels to be the same as
578        // the new label selector
579
580        Ok(())
581    }
582}
583
584#[async_trait::async_trait]
585impl k8s_controller::Context for Context {
586    type Resource = Console;
587    type Error = Error;
588
589    #[instrument(fields())]
590    async fn apply(
591        &self,
592        client: Client,
593        console: &Self::Resource,
594        metadata: &mut TraceMetadata,
595    ) -> Result<Option<Action>, Self::Error> {
596        if console.status.is_none() {
597            let step = metadata.step("initialize_status");
598            let console_api: Api<Console> =
599                Api::namespaced(client.clone(), &console.meta().namespace.clone().unwrap());
600            let mut new_console = console.clone();
601            new_console.status = Some(console.status());
602            console_api
603                .replace_status(
604                    &console.name_unchecked(),
605                    &PostParams::default(),
606                    &new_console,
607                )
608                .await?;
609            step.finish(Outcome::Completed);
610            // Updating the status should trigger a reconciliation
611            // which will include a status this time.
612            return Ok(None);
613        }
614
615        let namespace = console.namespace();
616        let network_policy_api: Api<NetworkPolicy> = Api::namespaced(client.clone(), &namespace);
617        let configmap_api: Api<ConfigMap> = Api::namespaced(client.clone(), &namespace);
618        let deployment_api: Api<Deployment> = Api::namespaced(client.clone(), &namespace);
619        let service_api: Api<Service> = Api::namespaced(client.clone(), &namespace);
620        let certificate_api: Api<Certificate> = Api::namespaced(client.clone(), &namespace);
621
622        let step = metadata.step("network_policies");
623        trace!("creating new network policies");
624        let network_policies = self.create_network_policies(console);
625        for network_policy in &network_policies {
626            apply_resource(&network_policy_api, network_policy).await?;
627        }
628        step.finish(if network_policies.is_empty() {
629            Outcome::Skipped
630        } else {
631            Outcome::Completed
632        });
633
634        let step = metadata.step("configmap");
635        trace!("creating new console configmap");
636        let console_configmap = self.create_console_app_configmap_object(console);
637        apply_resource(&configmap_api, &console_configmap).await?;
638        step.finish(Outcome::Completed);
639
640        let step = metadata.step("deployment");
641        trace!("creating new console deployment");
642        let console_deployment = self.create_console_deployment_object(console);
643        self.fix_deployment(&deployment_api, &console_deployment)
644            .await?;
645        apply_resource(&deployment_api, &console_deployment).await?;
646        step.finish(Outcome::Completed);
647
648        let step = metadata.step("service");
649        trace!("creating new console service");
650        let console_service = self.create_console_service_object(console);
651        apply_resource(&service_api, &console_service).await?;
652        step.finish(Outcome::Completed);
653
654        let step = metadata.step("certificate");
655        let console_external_certificate = self.create_console_external_certificate(console)?;
656        if let Some(certificate) = &console_external_certificate {
657            trace!("creating new console external certificate");
658            apply_resource(&certificate_api, certificate).await?;
659            step.finish(Outcome::Completed);
660        } else {
661            step.finish(Outcome::Skipped);
662        }
663
664        let step = metadata.step("sync_status");
665        self.sync_deployment_status(&client, console).await?;
666        step.finish(Outcome::Completed);
667
668        Ok(None)
669    }
670}