1use 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
54pub 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 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 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 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 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 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 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 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 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 }
572 Err(e) => return Err(anyhow::anyhow!(e).into()),
573 }
574
575 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 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}