Skip to main content

mz_orchestratord/controller/
materialize.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 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
60/// The name identifying this controller in its reconciliation metrics and in
61/// the reporter of the events it publishes.
62pub 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    /// The name of a `PriorityClass` to assign to environmentd pods, if any.
100    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    /// The name of a `PriorityClass` to assign to clusterd pods, if any.
105    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    /// Writes `status`, if it differs from the current status in anything but
177    /// its timestamps, and publishes the new `UpToDate` condition as an event
178    /// once it has been written.
179    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    /// Deletes `generation`'s resources, releasing the read holds its
220    /// environmentd was keeping.
221    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            // this is always a bug - we delay doing this check until the
294            // resource should have an environment id set, either from the
295            // license key, or explicitly given, or randomly defaulted.
296            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            // Updating the status should trigger a reconciliation
349            // which will include a status this time.
350            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                        // this makes it easier to use a license key in
389                        // development with no environment id set
390                        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            // Updating the spec should also trigger a reconciliation.
411            // We can't do that as part of the above check because you can't
412            // update both the spec and the status in a single api call.
413            return Ok(None);
414        }
415
416        if let Some(environment_id) = license_key_environment_id {
417            // we still allow a nil environment id in the license key to be
418            // accepted for any provided environment id, to support cloud
419            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        // we compare the hash against the environment resources generated
439        // for the current active generation, since that's what we expect to
440        // have been applied earlier, but we don't want to use these
441        // environment resources because when we apply them, we want to apply
442        // them with data that uses the new generation
443        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        // here we regenerate the environment resources using the
455        // same inputs except with an updated generation
456        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            // If we're in status promoting, we MUST promote now.
465            // We don't know if we successfully promoted or not yet.
466            (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            // There are changes pending, and we want to apply them.
479            (false, true, true) => {
480                // If a rollout has been in progress for longer than the
481                // configured timeout, cancel it. While a rollout is in
482                // progress the new generation runs un-promoted and holds back
483                // compaction via read holds; promoting it after a long delay
484                // can cause incident-inducing load, so we abort instead and
485                // let the user retry by requesting a fresh rollout.
486                //
487                // We never cancel a force-promoting rollout (including the
488                // `ImmediatelyPromoteCausingDowntime` strategy), because by
489                // then the previously-active generation may already be torn
490                // down, leaving nothing to fall back to.
491                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                            // Tear down the un-promoted generation to release
502                            // its read holds.
503                            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                                    // Mark this rollout request as completed so
518                                    // that we don't immediately retry it; the
519                                    // user must request a new rollout to try
520                                    // again.
521                                    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                // we remove the environment resources hash annotation here
593                // because if we fail halfway through applying the resources,
594                // things will be in an inconsistent state, and we don't want
595                // to allow the possibility of the user making a second
596                // change which reverts to the original state and then
597                // skipping retrying this apply, since that would leave
598                // things in a permanently inconsistent state.
599                // note that environment.spec will be empty here after
600                // replace_status, but this is fine because we already
601                // extracted all of the information we want from the spec
602                // earlier.
603                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                                // don't update the reconciliation id yet,
614                                // because the rollout hasn't yet completed. if
615                                // we fail later on, we want to ensure that the
616                                // rollout gets retried.
617                                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                    // The only reason someone would choose this strategy is if they didn't have
645                    // space for the two generations of pods.
646                    // Lets make room for the new ones by deleting the old generation.
647                    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                                        // Carry the `Applying` phase's
686                                        // timestamp forward (both phases are
687                                        // `Unknown`) so the rollout timeout
688                                        // spans Applying + ReadyToPromote
689                                        // rather than resetting here.
690                                        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                        // do this last, so that we keep traffic pointing at
707                        // the previous environmentd until the new one is
708                        // fully ready
709
710                        // Update the status before calling promote, so that we know
711                        // we've crossed the point of no return.
712                        // Once we see this status, we must promote without taking other actions.
713                        self.update_status(
714                            metadata,
715                            &mz_api,
716                            mz,
717                            MaterializeStatus {
718                                active_generation,
719                                // don't update the reconciliation id yet,
720                                // because the rollout hasn't yet completed. if
721                                // we fail later on, we want to ensure that the
722                                // rollout gets retried.
723                                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                                // also don't update the reconciliation id
763                                // here, because there was an error during
764                                // the rollout and we want to ensure it gets
765                                // retried.
766                                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            // There are changes pending, but we don't want to apply them yet.
793            (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            // No changes pending, but we might need to clean up a partially applied rollout.
830            (false, false, _) => {
831                // this can happen if we update the environment, but then revert
832                // that update before the update was deployed. in this case, we
833                // don't want the environment to still show up as
834                // WaitingForApproval.
835                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        // balancers rely on the environmentd service existing, which is
879        // enforced by the environmentd rollout process being able to call
880        // into the promotion endpoint
881
882        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        // and the console relies on the balancer service existing, which is
926        // enforced by wait_for_balancer
927
928        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
997/// The type of the event reporting that the `UpToDate` condition became
998/// `condition`: a warning if the environment is not up to date for any reason
999/// other than waiting for a rollout to be approved, which is the operator
1000/// following its configuration.
1001fn 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}