1use std::{
11 collections::BTreeSet,
12 sync::{Arc, Mutex},
13 time::Duration,
14};
15
16use anyhow::Context as _;
17use http::HeaderValue;
18use k8s_controller::{
19 Outcome, TraceMetadata,
20 events::{Event, EventType},
21};
22use k8s_openapi::{
23 api::core::v1::{Affinity, ResourceRequirements, Secret, Toleration},
24 apimachinery::pkg::apis::meta::v1::{Condition, Time},
25 jiff::{SignedDuration, Timestamp},
26};
27use kube::{
28 Api, Client, Resource, ResourceExt,
29 api::{ListParams, PostParams},
30 runtime::controller::Action,
31};
32use tracing::{debug, trace, warn};
33use uuid::Uuid;
34
35use crate::{
36 Error,
37 controller::materialize::generation::V161,
38 k8s::{apply_resource, delete_resource},
39 matching_image_from_environmentd_image_ref,
40 metrics::Metrics,
41 parse_image_tag,
42 tls::{DefaultCertificateSpecs, issuer_ref_defined, resolved_dns_names},
43};
44use mz_cloud_provider::CloudProvider;
45use mz_cloud_resources::crd::{
46 ManagedResource,
47 balancer::v1alpha1::{Balancer, BalancerSpec},
48 console::v1alpha1::{BalancerdRef, Console, ConsoleSpec, HttpConnectionScheme},
49 materialize::MaterializeRolloutStrategy,
50 materialize::v1alpha1::{Materialize, MaterializeStatus},
51};
52use mz_license_keys::validate;
53use mz_orchestrator_kubernetes::KubernetesImagePullPolicy;
54use mz_orchestrator_tracing::TracingCliArgs;
55use mz_ore::{cast::CastFrom, cli::KeyValueArg, instrument};
56
57pub mod generation;
58pub mod global;
59
60pub const CONTROLLER_NAME: &str = "materialize";
63
64#[derive(Clone)]
65pub struct Config {
66 pub cloud_provider: CloudProvider,
67 pub region: String,
68 pub create_balancers: bool,
69 pub create_console: bool,
70 pub helm_chart_version: Option<String>,
71 pub secrets_controller: String,
72 pub collect_pod_metrics: bool,
73 pub enable_prometheus_scrape_annotations: bool,
74
75 pub segment_api_key: Option<String>,
76 pub segment_client_side: bool,
77
78 pub console_image_tag_default: String,
79 pub console_image_tag_map: Vec<KeyValueArg<String, String>>,
80
81 pub aws_account_id: Option<String>,
82 pub environmentd_iam_role_arn: Option<String>,
83 pub environmentd_connection_role_arn: Option<String>,
84 pub aws_secrets_controller_tags: Vec<String>,
85 pub environmentd_availability_zones: Option<Vec<String>>,
86
87 pub ephemeral_volume_class: Option<String>,
88 pub scheduler_name: Option<String>,
89 pub enable_security_context: bool,
90 pub enable_internal_statement_logging: bool,
91 pub statement_logging_max_sample_rate: Option<f64>,
92 pub statement_logging_target_data_rate: Option<usize>,
93
94 pub orchestratord_pod_selector_labels: Vec<KeyValueArg<String, String>>,
95 pub environmentd_node_selector: Vec<KeyValueArg<String, String>>,
96 pub environmentd_affinity: Option<Affinity>,
97 pub environmentd_tolerations: Option<Vec<Toleration>>,
98 pub environmentd_default_resources: Option<ResourceRequirements>,
99 pub environmentd_priority_class_name: Option<String>,
101 pub clusterd_node_selector: Vec<KeyValueArg<String, String>>,
102 pub clusterd_affinity: Option<Affinity>,
103 pub clusterd_tolerations: Option<Vec<Toleration>>,
104 pub clusterd_priority_class_name: Option<String>,
106 pub image_pull_policy: KubernetesImagePullPolicy,
107 pub network_policies_internal_enabled: bool,
108 pub network_policies_ingress_enabled: bool,
109 pub network_policies_ingress_cidrs: Vec<String>,
110 pub network_policies_egress_enabled: bool,
111 pub network_policies_egress_cidrs: Vec<String>,
112
113 pub environmentd_cluster_replica_sizes: Option<String>,
114 pub bootstrap_default_cluster_replica_size: Option<String>,
115 pub bootstrap_builtin_system_cluster_replica_size: Option<String>,
116 pub bootstrap_builtin_probe_cluster_replica_size: Option<String>,
117 pub bootstrap_builtin_support_cluster_replica_size: Option<String>,
118 pub bootstrap_builtin_catalog_server_cluster_replica_size: Option<String>,
119 pub bootstrap_builtin_analytics_cluster_replica_size: Option<String>,
120 pub bootstrap_builtin_system_cluster_replication_factor: Option<u32>,
121 pub bootstrap_builtin_probe_cluster_replication_factor: Option<u32>,
122 pub bootstrap_builtin_support_cluster_replication_factor: Option<u32>,
123 pub bootstrap_builtin_analytics_cluster_replication_factor: Option<u32>,
124
125 pub environmentd_allowed_origins: Vec<HeaderValue>,
126 pub internal_console_proxy_url: String,
127
128 pub environmentd_sql_port: u16,
129 pub environmentd_http_port: u16,
130 pub environmentd_internal_sql_port: u16,
131 pub environmentd_internal_http_port: u16,
132 pub environmentd_internal_persist_pubsub_port: u16,
133
134 pub default_certificate_specs: DefaultCertificateSpecs,
135
136 pub disable_license_key_checks: bool,
137
138 pub tracing: TracingCliArgs,
139 pub orchestratord_namespace: String,
140}
141
142pub struct Context {
143 config: Config,
144 metrics: Arc<Metrics>,
145 needs_update: Arc<Mutex<BTreeSet<String>>>,
146}
147
148impl Context {
149 pub fn new(config: Config, metrics: Arc<Metrics>) -> Self {
150 if config.cloud_provider == CloudProvider::Aws {
151 assert!(
152 config.aws_account_id.is_some(),
153 "--aws-account-id is required when using --cloud-provider=aws"
154 );
155 }
156
157 Self {
158 config,
159 metrics,
160 needs_update: Default::default(),
161 }
162 }
163
164 fn set_needs_update(&self, mz: &Materialize, needs_update: bool) {
165 let mut needs_update_set = self.needs_update.lock().unwrap();
166 if needs_update {
167 needs_update_set.insert(mz.name_unchecked());
168 } else {
169 needs_update_set.remove(&mz.name_unchecked());
170 }
171 self.metrics
172 .environmentd_needs_update
173 .set(u64::cast_from(needs_update_set.len()));
174 }
175
176 async fn update_status(
180 &self,
181 metadata: &TraceMetadata,
182 mz_api: &Api<Materialize>,
183 mz: &Materialize,
184 status: MaterializeStatus,
185 needs_update: bool,
186 ) -> Result<Materialize, kube::Error> {
187 self.set_needs_update(mz, needs_update);
188
189 let mut new_mz = mz.clone();
190 if !mz
191 .status
192 .as_ref()
193 .map_or(true, |mz_status| mz_status.needs_update(&status))
194 {
195 return Ok(new_mz);
196 }
197
198 let condition = status.conditions.first().cloned();
199 new_mz.status = Some(status);
200 let new_mz = mz_api
201 .replace_status(&mz.name_unchecked(), &PostParams::default(), &new_mz)
202 .await?;
203
204 if let Some(condition) = condition {
205 metadata
206 .publish_event(Event {
207 type_: transition_event_type(&condition),
208 reason: condition.reason,
209 action: "Reconcile".into(),
210 note: Some(condition.message),
211 related: None,
212 })
213 .await;
214 }
215
216 Ok(new_mz)
217 }
218
219 async fn teardown_generation(
222 &self,
223 metadata: &TraceMetadata,
224 client: &Client,
225 mz: &Materialize,
226 resources: &generation::Resources,
227 generation: u64,
228 ) -> Result<(), Error> {
229 let step = metadata.step("teardown_generation");
230 resources
231 .teardown_generation(client, mz, generation)
232 .await?;
233 step.finish(Outcome::Completed);
234 Ok(())
235 }
236
237 async fn promote(
238 &self,
239 metadata: &TraceMetadata,
240 client: &Client,
241 mz: &Materialize,
242 resources: generation::Resources,
243 active_generation: u64,
244 desired_generation: u64,
245 resources_hash: String,
246 ) -> Result<Option<Action>, Error> {
247 let step = metadata.step("promote");
248 if let Some(action) = resources.promote_services(client, &mz.namespace()).await? {
249 step.finish(Outcome::Waiting);
250 return Ok(Some(action));
251 }
252 step.finish(Outcome::Completed);
253
254 self.teardown_generation(metadata, client, mz, &resources, active_generation)
255 .await?;
256 let mz_api: Api<Materialize> = Api::namespaced(client.clone(), &mz.namespace());
257 self.update_status(
258 metadata,
259 &mz_api,
260 mz,
261 MaterializeStatus {
262 active_generation: desired_generation,
263 last_completed_rollout_request: mz.requested_reconciliation_id(),
264 last_completed_rollout_environmentd_image_ref: Some(
265 mz.spec.environmentd_image_ref.clone(),
266 ),
267 resource_id: mz.status().resource_id,
268 resources_hash,
269 last_completed_rollout_hash: None,
270 conditions: vec![Condition {
271 type_: "UpToDate".into(),
272 status: "True".into(),
273 last_transition_time: Time(Timestamp::now()),
274 message: format!(
275 "Successfully applied changes for generation {desired_generation}"
276 ),
277 observed_generation: mz.meta().generation,
278 reason: "Applied".into(),
279 }],
280 },
281 false,
282 )
283 .await?;
284 Ok(None)
285 }
286
287 async fn check_environment_id_conflicts(
288 &self,
289 client: &Client,
290 mz: &Materialize,
291 ) -> Result<(), Error> {
292 if mz.spec.environment_id.is_nil() {
293 return Err(Error::Anyhow(anyhow::anyhow!(
297 "trying to reconcile a materialize resource with no environment id - this is a bug!"
298 )));
299 }
300
301 let mz_api: Api<Materialize> = Api::all(client.clone());
302 let all_mz = mz_api.list(&ListParams::default()).await?;
303 for existing_mz in &all_mz.items {
304 if existing_mz.spec.environment_id == mz.spec.environment_id
305 && existing_mz.metadata.uid != mz.metadata.uid
306 {
307 return Err(Error::Anyhow(anyhow::anyhow!(
308 "Materialize resources {}/{} and {}/{} have the environmentId field set to the same value. This field must be unique across environments.",
309 mz.namespace(),
310 mz.name_unchecked(),
311 existing_mz.namespace(),
312 existing_mz.name_unchecked(),
313 )));
314 }
315 }
316
317 Ok(())
318 }
319}
320
321#[async_trait::async_trait]
322impl k8s_controller::Context for Context {
323 type Resource = Materialize;
324 type Error = Error;
325
326 const FINALIZER_NAME: Option<&'static str> =
327 Some("orchestratord.materialize.cloud/materialize");
328
329 #[instrument(fields(organization_name=mz.name_unchecked()))]
330 async fn apply(
331 &self,
332 client: Client,
333 mz: &Self::Resource,
334 metadata: &mut TraceMetadata,
335 ) -> Result<Option<Action>, Self::Error> {
336 let metadata = &*metadata;
337 let mz_api: Api<Materialize> = Api::namespaced(client.clone(), &mz.namespace());
338 let balancer_api: Api<Balancer> = Api::namespaced(client.clone(), &mz.namespace());
339 let console_api: Api<Console> = Api::namespaced(client.clone(), &mz.namespace());
340 let secret_api: Api<Secret> = Api::namespaced(client.clone(), &mz.namespace());
341
342 let status = mz.status();
343 if mz.status.is_none() {
344 let step = metadata.step("initialize_status");
345 self.update_status(metadata, &mz_api, mz, status, true)
346 .await?;
347 step.finish(Outcome::Completed);
348 return Ok(None);
351 }
352
353 let step = metadata.step("resolve_environment_id");
354 let backend_secret = secret_api.get(&mz.spec.backend_secret_name).await?;
355 let license_key_environment_id: Option<Uuid> = if let Some(license_key) = backend_secret
356 .data
357 .as_ref()
358 .and_then(|data| data.get("license_key"))
359 {
360 let license_key = validate(
361 str::from_utf8(&license_key.0)
362 .context("invalid utf8")?
363 .trim(),
364 )?;
365 let environment_id = license_key
366 .environment_id
367 .parse()
368 .context("invalid environment id in license key")?;
369 Some(environment_id)
370 } else {
371 if mz.meets_minimum_version(&V161) {
372 return Err(Error::Anyhow(anyhow::anyhow!(
373 "license_key is required when running in kubernetes",
374 )));
375 } else {
376 None
377 }
378 };
379
380 if mz.spec.request_rollout.is_nil() || mz.spec.environment_id.is_nil() {
381 let mut mz = mz.clone();
382 if mz.spec.request_rollout.is_nil() {
383 mz.spec.request_rollout = Uuid::new_v4();
384 }
385 if mz.spec.environment_id.is_nil() {
386 if let Some(environment_id) = license_key_environment_id {
387 if environment_id.is_nil() {
388 mz.spec.environment_id = Uuid::new_v4();
391 } else {
392 mz.spec.environment_id = environment_id;
393 }
394 } else {
395 if mz.meets_minimum_version(&V161) {
396 return Err(Error::Anyhow(anyhow::anyhow!(
397 "environmentId is not set in materialize resource {}/{} but no license key was given",
398 mz.namespace(),
399 mz.name_unchecked()
400 )));
401 } else {
402 mz.spec.environment_id = Uuid::new_v4();
403 }
404 }
405 }
406 mz_api
407 .replace(&mz.name_unchecked(), &PostParams::default(), &mz)
408 .await?;
409 step.finish(Outcome::Completed);
410 return Ok(None);
414 }
415
416 if let Some(environment_id) = license_key_environment_id {
417 if !environment_id.is_nil() && mz.spec.environment_id != environment_id {
420 return Err(Error::Anyhow(anyhow::anyhow!(
421 "environment_id is set in materialize resource {}/{} but does not match the environment_id set in the associated license key {}",
422 mz.namespace(),
423 mz.name_unchecked(),
424 environment_id,
425 )));
426 }
427 }
428
429 self.check_environment_id_conflicts(&client, mz).await?;
430 step.finish(Outcome::Completed);
431
432 let step = metadata.step("global_resources");
433 global::Resources::new(&self.config, mz)?
434 .apply(&client, &mz.namespace())
435 .await?;
436 step.finish(Outcome::Completed);
437
438 let active_resources =
444 generation::Resources::new(&self.config, mz, status.active_generation);
445 let has_current_changes = status.resources_hash != active_resources.generate_hash();
446 let active_generation = status.active_generation;
447 let next_generation = active_generation + 1;
448 let desired_generation = if has_current_changes {
449 next_generation
450 } else {
451 active_generation
452 };
453
454 let resources = generation::Resources::new(&self.config, mz, desired_generation);
457 let resources_hash = resources.generate_hash();
458
459 let mut result = match (
460 mz.is_promoting(),
461 has_current_changes,
462 mz.rollout_requested(),
463 ) {
464 (true, _, _) => {
467 self.promote(
468 metadata,
469 &client,
470 mz,
471 resources,
472 active_generation,
473 desired_generation,
474 resources_hash,
475 )
476 .await
477 }
478 (false, true, true) => {
480 if !mz.should_force_promote() {
492 if let Some(started) = mz.rollout_in_progress_since() {
493 let timeout = mz.rollout_request_timeout();
494 let elapsed = Timestamp::now().duration_since(started);
495 let timed_out = SignedDuration::try_from(timeout)
496 .is_ok_and(|timeout| elapsed >= timeout);
497 if timed_out {
498 warn!(
499 "rollout to generation {desired_generation} exceeded timeout, cancelling"
500 );
501 self.teardown_generation(
504 metadata,
505 &client,
506 mz,
507 &resources,
508 next_generation,
509 )
510 .await?;
511 self.update_status(
512 metadata,
513 &mz_api,
514 mz,
515 MaterializeStatus {
516 active_generation,
517 last_completed_rollout_request: mz
522 .requested_reconciliation_id(),
523 last_completed_rollout_environmentd_image_ref: status
524 .last_completed_rollout_environmentd_image_ref
525 .clone(),
526 resource_id: status.resource_id.clone(),
527 resources_hash: status.resources_hash.clone(),
528 last_completed_rollout_hash: None,
529 conditions: vec![Condition {
530 type_: "UpToDate".into(),
531 status: "False".into(),
532 last_transition_time: Time(Timestamp::now()),
533 message: format!(
534 "Cancelled rollout to generation \
535 {desired_generation} after it \
536 exceeded the rollout timeout of {}",
537 humantime::format_duration(timeout),
538 ),
539 observed_generation: mz.meta().generation,
540 reason: "RolloutTimeout".into(),
541 }],
542 },
543 active_generation != desired_generation,
544 )
545 .await?;
546 return Ok(None);
547 }
548 }
549 }
550
551 if !mz.within_upgrade_window() {
552 let last_completed_rollout_environmentd_image_ref =
553 status.last_completed_rollout_environmentd_image_ref;
554
555 self.update_status(
556 metadata,
557 &mz_api,
558 mz,
559 MaterializeStatus {
560 active_generation,
561 last_completed_rollout_request: status.last_completed_rollout_request,
562 last_completed_rollout_environmentd_image_ref:
563 last_completed_rollout_environmentd_image_ref.clone(),
564 resource_id: status.resource_id,
565 resources_hash: status.resources_hash,
566 last_completed_rollout_hash: None,
567 conditions: vec![Condition {
568 type_: "UpToDate".into(),
569 status: "False".into(),
570 last_transition_time: Time(Timestamp::now()),
571 message: format!(
572 "Refusing to upgrade from {} to {}. \
573 More than one major version from \
574 last successful rollout. If coming \
575 from Self Managed 25.2, upgrade to \
576 materialize/environmentd:v0.147.20 \
577 first.",
578 last_completed_rollout_environmentd_image_ref
579 .expect("should be set if upgrade window check fails"),
580 mz.spec.environmentd_image_ref,
581 ),
582 observed_generation: mz.meta().generation,
583 reason: "FailedDeploy".into(),
584 }],
585 },
586 active_generation != desired_generation,
587 )
588 .await?;
589 return Ok(None);
590 }
591
592 let mz = if mz.is_ready_to_promote(&resources_hash) {
604 mz
605 } else {
606 &self
607 .update_status(
608 metadata,
609 &mz_api,
610 mz,
611 MaterializeStatus {
612 active_generation,
613 last_completed_rollout_request: status
618 .last_completed_rollout_request,
619 last_completed_rollout_environmentd_image_ref: status
620 .last_completed_rollout_environmentd_image_ref,
621 resource_id: status.resource_id.clone(),
622 resources_hash: String::new(),
623 last_completed_rollout_hash: None,
624 conditions: vec![Condition {
625 type_: "UpToDate".into(),
626 status: "Unknown".into(),
627 last_transition_time: Time(Timestamp::now()),
628 message: format!(
629 "Applying changes for generation {desired_generation}"
630 ),
631 observed_generation: mz.meta().generation,
632 reason: "Applying".into(),
633 }],
634 },
635 active_generation != desired_generation,
636 )
637 .await?
638 };
639 let status = mz.status();
640
641 if mz.spec.rollout_strategy
642 == MaterializeRolloutStrategy::ImmediatelyPromoteCausingDowntime
643 {
644 self.teardown_generation(metadata, &client, mz, &resources, active_generation)
648 .await?;
649 }
650
651 trace!("applying environment resources");
652 let step = metadata.step("generation_resources");
653 let applied = resources
654 .apply(&client, mz.should_force_promote(), &mz.namespace())
655 .await;
656 step.finish_with(&applied);
657 match applied {
658 Ok(Some(action)) => {
659 trace!("new environment is not yet ready");
660 Ok(Some(action))
661 }
662 Ok(None) => {
663 if mz.spec.rollout_strategy == MaterializeRolloutStrategy::ManuallyPromote
664 && !mz.should_force_promote()
665 {
666 trace!(
667 "Ready to promote, but not promoting because the instance is configured with ManuallyPromote rollout strategy."
668 );
669 self.update_status(
670 metadata,
671 &mz_api,
672 mz,
673 MaterializeStatus {
674 active_generation,
675 last_completed_rollout_request: status
676 .last_completed_rollout_request,
677 last_completed_rollout_environmentd_image_ref: status
678 .last_completed_rollout_environmentd_image_ref,
679 resource_id: status.resource_id,
680 resources_hash,
681 last_completed_rollout_hash: None,
682 conditions: vec![Condition {
683 type_: "UpToDate".into(),
684 status: "Unknown".into(),
685 last_transition_time: Time(mz.up_to_date_transition_time(
691 "Unknown",
692 Timestamp::now(),
693 )),
694 message: format!(
695 "Ready to promote generation {desired_generation}"
696 ),
697 observed_generation: mz.meta().generation,
698 reason: "ReadyToPromote".into(),
699 }],
700 },
701 active_generation != desired_generation,
702 )
703 .await?;
704 return Ok(None);
705 }
706 self.update_status(
714 metadata,
715 &mz_api,
716 mz,
717 MaterializeStatus {
718 active_generation,
719 last_completed_rollout_request: status
724 .last_completed_rollout_request,
725 last_completed_rollout_environmentd_image_ref: status
726 .last_completed_rollout_environmentd_image_ref,
727 resource_id: status.resource_id,
728 resources_hash: resources_hash.clone(),
729 last_completed_rollout_hash: None,
730 conditions: vec![Condition {
731 type_: "UpToDate".into(),
732 status: "Unknown".into(),
733 last_transition_time: Time(Timestamp::now()),
734 message: format!(
735 "Attempting to promote generation {desired_generation}"
736 ),
737 observed_generation: mz.meta().generation,
738 reason: "Promoting".into(),
739 }],
740 },
741 active_generation != desired_generation,
742 )
743 .await?;
744 self.promote(
745 metadata,
746 &client,
747 mz,
748 resources,
749 active_generation,
750 desired_generation,
751 resources_hash,
752 )
753 .await
754 }
755 Err(e) => {
756 self.update_status(
757 metadata,
758 &mz_api,
759 mz,
760 MaterializeStatus {
761 active_generation,
762 last_completed_rollout_request: status
767 .last_completed_rollout_request,
768 last_completed_rollout_environmentd_image_ref: status
769 .last_completed_rollout_environmentd_image_ref,
770 resource_id: status.resource_id,
771 resources_hash: status.resources_hash,
772 last_completed_rollout_hash: None,
773 conditions: vec![Condition {
774 type_: "UpToDate".into(),
775 status: "False".into(),
776 last_transition_time: Time(Timestamp::now()),
777 message: format!(
778 "Failed to apply changes for \
779 generation {desired_generation}: {e}"
780 ),
781 observed_generation: mz.meta().generation,
782 reason: "FailedDeploy".into(),
783 }],
784 },
785 active_generation != desired_generation,
786 )
787 .await?;
788 Err(e)
789 }
790 }
791 }
792 (false, true, false) => {
794 let mut needs_update = mz.conditions_need_update();
795 if mz.update_in_progress() {
796 self.teardown_generation(metadata, &client, mz, &resources, next_generation)
797 .await?;
798 needs_update = true;
799 }
800 if needs_update {
801 self.update_status(metadata, &mz_api,
802 mz,
803 MaterializeStatus {
804 active_generation,
805 last_completed_rollout_request: mz.requested_reconciliation_id(),
806 last_completed_rollout_environmentd_image_ref: status
807 .last_completed_rollout_environmentd_image_ref,
808 resource_id: status.resource_id.clone(),
809 resources_hash: status.resources_hash,
810 last_completed_rollout_hash: None,
811 conditions: vec![Condition {
812 type_: "UpToDate".into(),
813 status: "False".into(),
814 last_transition_time: Time(Timestamp::now()),
815 message: format!(
816 "Changes detected, waiting for approval for generation {desired_generation}"
817 ),
818 observed_generation: mz.meta().generation,
819 reason: "WaitingForApproval".into(),
820 }],
821 },
822 active_generation != desired_generation,
823 )
824 .await?;
825 }
826 debug!("changes detected, waiting for approval");
827 Ok(None)
828 }
829 (false, false, _) => {
831 let mut needs_update = mz.conditions_need_update() || mz.rollout_requested();
836 if mz.update_in_progress() {
837 self.teardown_generation(metadata, &client, mz, &resources, next_generation)
838 .await?;
839 needs_update = true;
840 }
841 if needs_update {
842 self.update_status(
843 metadata,
844 &mz_api,
845 mz,
846 MaterializeStatus {
847 active_generation,
848 last_completed_rollout_request: mz.requested_reconciliation_id(),
849 last_completed_rollout_environmentd_image_ref: status
850 .last_completed_rollout_environmentd_image_ref,
851 resource_id: status.resource_id.clone(),
852 resources_hash: status.resources_hash,
853 last_completed_rollout_hash: None,
854 conditions: vec![Condition {
855 type_: "UpToDate".into(),
856 status: "True".into(),
857 last_transition_time: Time(Timestamp::now()),
858 message: format!(
859 "No changes found from generation {active_generation}"
860 ),
861 observed_generation: mz.meta().generation,
862 reason: "Applied".into(),
863 }],
864 },
865 active_generation != desired_generation,
866 )
867 .await?;
868 }
869 debug!("no changes");
870 Ok(None)
871 }
872 }?;
873
874 if let Some(action) = result {
875 return Ok(Some(action));
876 }
877
878 let step = metadata.step("balancer");
883 if self.config.create_balancers {
884 let balancer = Balancer {
885 metadata: mz.managed_resource_meta(mz.name_unchecked()),
886 spec: BalancerSpec {
887 balancerd_image_ref: matching_image_from_environmentd_image_ref(
888 mz.active_environmentd_image_ref(),
889 "balancerd",
890 None,
891 ),
892 resource_requirements: mz.spec.balancerd_resource_requirements.clone(),
893 configmap_name: mz.spec.balancerd_configmap_name.clone(),
894 replicas: Some(mz.balancerd_replicas()),
895 external_certificate_spec: mz.spec.balancerd_external_certificate_spec.clone(),
896 internal_certificate_spec: mz.spec.internal_certificate_spec.clone(),
897 pod_annotations: mz.spec.pod_annotations.clone(),
898 pod_labels: mz.spec.pod_labels.clone(),
899 static_routing: Some(
900 mz_cloud_resources::crd::balancer::v1alpha1::StaticRoutingConfig {
901 environmentd_namespace: mz.namespace(),
902 environmentd_service_name: mz.environmentd_service_name(),
903 },
904 ),
905 frontegg_routing: None,
906 resource_id: Some(status.resource_id.clone()),
907 },
908 status: None,
909 };
910 let balancer = apply_resource(&balancer_api, &balancer).await?;
911 result = wait_for_balancer(&balancer)?;
912 step.finish(match result {
913 Some(_) => Outcome::Waiting,
914 None => Outcome::Completed,
915 });
916 } else {
917 delete_resource(&balancer_api, &mz.name_unchecked()).await?;
918 step.finish(Outcome::Skipped);
919 }
920
921 if let Some(action) = result {
922 return Ok(Some(action));
923 }
924
925 let step = metadata.step("console");
929 if self.config.create_console {
930 let active_environmentd_image_ref = mz.active_environmentd_image_ref();
931 let environmentd_image_tag =
932 parse_image_tag(active_environmentd_image_ref).unwrap_or("latest");
933 let console_image_tag = self
934 .config
935 .console_image_tag_map
936 .iter()
937 .find(|kv| kv.key == environmentd_image_tag)
938 .map(|kv| kv.value.clone())
939 .unwrap_or_else(|| self.config.console_image_tag_default.clone());
940 let console = Console {
941 metadata: mz.managed_resource_meta(mz.name_unchecked()),
942 spec: ConsoleSpec {
943 console_image_ref: matching_image_from_environmentd_image_ref(
944 active_environmentd_image_ref,
945 "console",
946 Some(&console_image_tag),
947 ),
948 resource_requirements: mz.spec.console_resource_requirements.clone(),
949 replicas: Some(mz.console_replicas()),
950 external_certificate_spec: mz.spec.console_external_certificate_spec.clone(),
951 pod_annotations: mz.spec.pod_annotations.clone(),
952 pod_labels: mz.spec.pod_labels.clone(),
953 balancerd: BalancerdRef {
954 service_name: mz.balancerd_service_name(),
955 namespace: mz.namespace(),
956 scheme: if issuer_ref_defined(
957 &self.config.default_certificate_specs.balancerd_external,
958 &mz.spec.balancerd_external_certificate_spec,
959 ) {
960 HttpConnectionScheme::Https
961 } else {
962 HttpConnectionScheme::Http
963 },
964 dns_names: resolved_dns_names(
965 &self.config.default_certificate_specs.balancerd_external,
966 &mz.spec.balancerd_external_certificate_spec,
967 ),
968 },
969 authenticator_kind: mz.spec.authenticator_kind,
970 resource_id: Some(status.resource_id),
971 },
972 status: None,
973 };
974 apply_resource(&console_api, &console).await?;
975 step.finish(Outcome::Completed);
976 } else {
977 delete_resource(&console_api, &mz.name_unchecked()).await?;
978 step.finish(Outcome::Skipped);
979 }
980
981 Ok(result)
982 }
983
984 #[instrument(fields(organization_name=mz.name_unchecked()))]
985 async fn cleanup(
986 &self,
987 _client: Client,
988 mz: &Self::Resource,
989 _metadata: &mut TraceMetadata,
990 ) -> Result<Option<Action>, Self::Error> {
991 self.set_needs_update(mz, false);
992
993 Ok(None)
994 }
995}
996
997fn transition_event_type(condition: &Condition) -> EventType {
1002 match condition.reason.as_str() {
1003 "WaitingForApproval" => EventType::Normal,
1004 _ if condition.status == "False" => EventType::Warning,
1005 _ => EventType::Normal,
1006 }
1007}
1008
1009fn wait_for_balancer(balancer: &Balancer) -> Result<Option<Action>, Error> {
1010 if let Some(conditions) = balancer
1011 .status
1012 .as_ref()
1013 .map(|status| status.conditions.as_slice())
1014 {
1015 if conditions
1016 .iter()
1017 .any(|condition| condition.type_ == "Ready" && condition.status == "True")
1018 {
1019 return Ok(None);
1020 }
1021 }
1022
1023 Ok(Some(Action::requeue(Duration::from_secs(1))))
1024}
1025
1026#[cfg(test)]
1027mod tests {
1028 use super::*;
1029
1030 fn condition(status: &str, reason: &str) -> Condition {
1031 Condition {
1032 type_: "UpToDate".into(),
1033 status: status.into(),
1034 reason: reason.into(),
1035 message: String::new(),
1036 last_transition_time: Time(Timestamp::now()),
1037 observed_generation: None,
1038 }
1039 }
1040
1041 #[mz_ore::test]
1042 fn test_transition_event_type() {
1043 for (status, reason, expected) in [
1044 ("True", "Applied", EventType::Normal),
1045 ("Unknown", "Applying", EventType::Normal),
1046 ("Unknown", "ReadyToPromote", EventType::Normal),
1047 ("Unknown", "Promoting", EventType::Normal),
1048 ("False", "WaitingForApproval", EventType::Normal),
1049 ("False", "FailedDeploy", EventType::Warning),
1050 ("False", "RolloutTimeout", EventType::Warning),
1051 ("False", "SomeFutureFailure", EventType::Warning),
1052 ("Unknown", "SomeFuturePhase", EventType::Normal),
1053 ] {
1054 assert_eq!(
1055 transition_event_type(&condition(status, reason)),
1056 expected,
1057 "status={status} reason={reason}",
1058 );
1059 }
1060 }
1061}