Skip to main content

mz_controller/
clusters.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
10//! Cluster management.
11
12use std::collections::{BTreeMap, BTreeSet};
13use std::fmt;
14use std::num::NonZero;
15use std::str::FromStr;
16use std::sync::Arc;
17use std::sync::LazyLock;
18use std::time::Duration;
19
20use anyhow::anyhow;
21use bytesize::ByteSize;
22use chrono::{DateTime, Utc};
23use futures::stream::{BoxStream, StreamExt};
24use mz_cluster_client::client::{ClusterReplicaLocation, TimelyConfig};
25use mz_compute_client::logging::LogVariant;
26use mz_compute_types::config::{ComputeReplicaConfig, ComputeReplicaLogging};
27use mz_controller_types::dyncfgs::{
28    ARRANGEMENT_EXERT_PROPORTIONALITY, CONTROLLER_PAST_GENERATION_REPLICA_CLEANUP_RETRY_INTERVAL,
29    ENABLE_TIMELY_ZERO_COPY, ENABLE_TIMELY_ZERO_COPY_LGALLOC, ENABLE_UNIFIED_CLUSTER,
30    TIMELY_ZERO_COPY_LIMIT,
31};
32use mz_controller_types::{ClusterId, ReplicaId};
33use mz_orchestrator::NamespacedOrchestrator;
34use mz_orchestrator::{
35    CpuLimit, DiskLimit, LabelSelectionLogic, LabelSelector, MemoryLimit, Service, ServiceConfig,
36    ServiceEvent, ServicePort,
37};
38use mz_ore::cast::CastInto;
39use mz_ore::task::{self, AbortOnDropHandle};
40use mz_ore::{halt, instrument};
41use mz_repr::GlobalId;
42use mz_repr::adt::numeric::Numeric;
43use regex::Regex;
44use serde::{Deserialize, Serialize};
45use tokio::time;
46use tracing::{error, info, warn};
47
48use crate::Controller;
49
50/// Configures a cluster.
51pub struct ClusterConfig {
52    /// The logging variants to enable on the compute instance.
53    ///
54    /// Each logging variant is mapped to the identifier under which to register
55    /// the arrangement storing the log's data.
56    pub arranged_logs: BTreeMap<LogVariant, GlobalId>,
57    /// An optional arbitrary string that describes the class of the workload
58    /// this cluster is running (e.g., `production` or `staging`).
59    pub workload_class: Option<String>,
60}
61
62/// The status of a cluster.
63pub type ClusterStatus = mz_orchestrator::ServiceStatus;
64
65/// Configures a cluster replica.
66#[derive(Clone, Debug, Serialize, PartialEq)]
67pub struct ReplicaConfig {
68    /// The location of the replica.
69    pub location: ReplicaLocation,
70    /// Configuration for the compute half of the replica.
71    pub compute: ComputeReplicaConfig,
72}
73
74/// Configures the resource allocation for a cluster replica.
75#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
76pub struct ReplicaAllocation {
77    /// The memory limit for each process in the replica.
78    pub memory_limit: Option<MemoryLimit>,
79    /// The CPU limit for each process in the replica.
80    pub cpu_limit: Option<CpuLimit>,
81    /// The CPU limit for each process in the replica.
82    pub cpu_request: Option<CpuLimit>,
83    /// The disk limit for each process in the replica.
84    pub disk_limit: Option<DiskLimit>,
85    /// The number of processes in the replica.
86    pub scale: NonZero<u16>,
87    /// The number of worker threads in the replica.
88    pub workers: NonZero<usize>,
89    /// The number of credits per hour that the replica consumes.
90    #[serde(deserialize_with = "mz_repr::adt::numeric::str_serde::deserialize")]
91    pub credits_per_hour: Numeric,
92    /// Whether each process has exclusive access to its CPU cores.
93    #[serde(default)]
94    pub cpu_exclusive: bool,
95    /// Whether this size represents a modern "cc" size rather than a legacy
96    /// T-shirt size.
97    #[serde(default = "default_true")]
98    pub is_cc: bool,
99    /// The size *family* this size belongs to, e.g. the size `D.1-xsmall`
100    /// belongs to family `D` and the legacy t-shirt sizes belong to family
101    /// `legacy`. The family is the coarse axis and is *not* a prefix of the size
102    /// name in general. Used as the
103    /// `replica_size_family` attribute when evaluating replica-local scoped
104    /// feature flags (see the scoped feature flags design). When unset, the
105    /// family falls back to a value derived from [`Self::is_cc`] via
106    /// [`ReplicaAllocation::family`].
107    #[serde(default)]
108    pub family: Option<String>,
109    /// Whether instances of this type use swap as the spill-to-disk mechanism.
110    #[serde(default)]
111    pub swap_enabled: bool,
112    /// Whether instances of this type can be created.
113    #[serde(default)]
114    pub disabled: bool,
115    /// Additional node selectors.
116    #[serde(default)]
117    pub selectors: BTreeMap<String, String>,
118}
119
120impl ReplicaAllocation {
121    /// The name of the size family this allocation belongs to, used as the
122    /// `replica_size_family` attribute when evaluating replica-local scoped
123    /// feature flags.
124    ///
125    /// Falls back to a value derived from [`Self::is_cc`] when [`Self::family`]
126    /// is unset: `"cc"` for modern sizes and `"legacy"` for the legacy t-shirt
127    /// sizes. This keeps the legacy family targetable even before every size
128    /// gains an explicit `family` in the size configuration.
129    pub fn family(&self) -> &str {
130        match &self.family {
131            Some(family) => family.as_str(),
132            None if self.is_cc => "cc",
133            None => "legacy",
134        }
135    }
136}
137
138fn default_true() -> bool {
139    true
140}
141
142#[mz_ore::test]
143// We test this particularly because we deserialize values from strings.
144#[cfg_attr(miri, ignore)] // unsupported operation: can't call foreign function `decContextDefault` on OS `linux`
145fn test_replica_allocation_deserialization() {
146    use bytesize::ByteSize;
147    use mz_ore::{assert_err, assert_ok};
148
149    let data = r#"
150        {
151            "cpu_limit": 1.0,
152            "memory_limit": "10GiB",
153            "disk_limit": "100MiB",
154            "scale": 16,
155            "workers": 1,
156            "credits_per_hour": "16",
157            "swap_enabled": true,
158            "selectors": {
159                "key1": "value1",
160                "key2": "value2"
161            }
162        }"#;
163
164    let replica_allocation: ReplicaAllocation = serde_json::from_str(data)
165        .expect("deserialization from JSON succeeds for ReplicaAllocation");
166
167    assert_eq!(
168        replica_allocation,
169        ReplicaAllocation {
170            credits_per_hour: 16.into(),
171            disk_limit: Some(DiskLimit(ByteSize::mib(100))),
172            disabled: false,
173            memory_limit: Some(MemoryLimit(ByteSize::gib(10))),
174            cpu_limit: Some(CpuLimit::from_millicpus(1000)),
175            cpu_request: None,
176            cpu_exclusive: false,
177            is_cc: true,
178            family: None,
179            swap_enabled: true,
180            scale: NonZero::new(16).unwrap(),
181            workers: NonZero::new(1).unwrap(),
182            selectors: BTreeMap::from([
183                ("key1".to_string(), "value1".to_string()),
184                ("key2".to_string(), "value2".to_string())
185            ]),
186        }
187    );
188
189    let data = r#"
190        {
191            "cpu_limit": 0,
192            "memory_limit": "0GiB",
193            "disk_limit": "0MiB",
194            "scale": 1,
195            "workers": 1,
196            "credits_per_hour": "0",
197            "cpu_exclusive": true,
198            "disabled": true
199        }"#;
200
201    let replica_allocation: ReplicaAllocation = serde_json::from_str(data)
202        .expect("deserialization from JSON succeeds for ReplicaAllocation");
203
204    assert_eq!(
205        replica_allocation,
206        ReplicaAllocation {
207            credits_per_hour: 0.into(),
208            disk_limit: Some(DiskLimit(ByteSize::mib(0))),
209            disabled: true,
210            memory_limit: Some(MemoryLimit(ByteSize::gib(0))),
211            cpu_limit: Some(CpuLimit::from_millicpus(0)),
212            cpu_request: None,
213            cpu_exclusive: true,
214            is_cc: true,
215            family: None,
216            swap_enabled: false,
217            scale: NonZero::new(1).unwrap(),
218            workers: NonZero::new(1).unwrap(),
219            selectors: Default::default(),
220        }
221    );
222
223    // `scale` and `workers` must be non-zero.
224    let data = r#"{"scale": 0, "workers": 1, "credits_per_hour": "0"}"#;
225    assert_err!(serde_json::from_str::<ReplicaAllocation>(data));
226    let data = r#"{"scale": 1, "workers": 0, "credits_per_hour": "0"}"#;
227    assert_err!(serde_json::from_str::<ReplicaAllocation>(data));
228    let data = r#"{"scale": 1, "workers": 1, "credits_per_hour": "0"}"#;
229    assert_ok!(serde_json::from_str::<ReplicaAllocation>(data));
230}
231
232#[mz_ore::test]
233#[cfg_attr(miri, ignore)] // unsupported operation: can't call foreign function `decContextDefault` on OS `linux`
234fn test_replica_allocation_family() {
235    let parse = |json: &str| -> ReplicaAllocation {
236        serde_json::from_str(json).expect("deserialization from JSON succeeds")
237    };
238
239    // An explicit `family` is used verbatim.
240    assert_eq!(
241        parse(r#"{"scale": 1, "workers": 1, "credits_per_hour": "0", "family": "D"}"#).family(),
242        "D"
243    );
244    // Without an explicit `family`, modern (`is_cc`) sizes fall back to "cc".
245    // `is_cc` defaults to true.
246    assert_eq!(
247        parse(r#"{"scale": 1, "workers": 1, "credits_per_hour": "0"}"#).family(),
248        "cc"
249    );
250    // Without an explicit `family`, legacy (non-`is_cc`) sizes fall back to
251    // "legacy".
252    assert_eq!(
253        parse(r#"{"scale": 1, "workers": 1, "credits_per_hour": "0", "is_cc": false}"#).family(),
254        "legacy"
255    );
256    // An explicit family wins even for a legacy size.
257    assert_eq!(
258        parse(
259            r#"{"scale": 1, "workers": 1, "credits_per_hour": "0", "is_cc": false, "family": "legacy-special"}"#
260        )
261        .family(),
262        "legacy-special"
263    );
264}
265
266/// Configures the location of a cluster replica.
267#[derive(Clone, Debug, Serialize, PartialEq)]
268pub enum ReplicaLocation {
269    /// An unmanaged replica.
270    Unmanaged(UnmanagedReplicaLocation),
271    /// A managed replica.
272    Managed(ManagedReplicaLocation),
273}
274
275impl ReplicaLocation {
276    /// Returns the number of processes specified by this replica location.
277    pub fn num_processes(&self) -> usize {
278        match self {
279            ReplicaLocation::Unmanaged(UnmanagedReplicaLocation {
280                computectl_addrs, ..
281            }) => computectl_addrs.len(),
282            ReplicaLocation::Managed(ManagedReplicaLocation { allocation, .. }) => {
283                allocation.scale.cast_into()
284            }
285        }
286    }
287
288    pub fn billed_as(&self) -> Option<&str> {
289        match self {
290            ReplicaLocation::Managed(ManagedReplicaLocation { billed_as, .. }) => {
291                billed_as.as_deref()
292            }
293            ReplicaLocation::Unmanaged(_) => None,
294        }
295    }
296
297    pub fn internal(&self) -> bool {
298        match self {
299            ReplicaLocation::Managed(ManagedReplicaLocation { internal, .. }) => *internal,
300            ReplicaLocation::Unmanaged(_) => false,
301        }
302    }
303
304    /// Returns the number of workers specified by this replica location.
305    ///
306    /// `None` for unmanaged replicas, whose worker count we don't know.
307    pub fn workers(&self) -> Option<usize> {
308        match self {
309            ReplicaLocation::Managed(ManagedReplicaLocation { allocation, .. }) => {
310                Some(allocation.workers.get() * self.num_processes())
311            }
312            ReplicaLocation::Unmanaged(_) => None,
313        }
314    }
315
316    /// Whether the replica is durably marked `pending`.
317    ///
318    /// Vestigial: no path creates one anymore. A crash on a version that still
319    /// staged reconfigurations through overlap replicas could have left one
320    /// behind, and the catalog-open migration reaps those.
321    pub fn pending(&self) -> bool {
322        match self {
323            ReplicaLocation::Managed(ManagedReplicaLocation { pending, .. }) => *pending,
324            ReplicaLocation::Unmanaged(_) => false,
325        }
326    }
327}
328
329/// The "role" of a cluster, which is currently used to determine the
330/// severity of alerts for problems with its replicas.
331#[derive(Debug, Clone)]
332pub enum ClusterRole {
333    /// The existence and proper functioning of the cluster's replicas is
334    /// business-critical for Materialize.
335    SystemCritical,
336    /// Assuming no bugs, the cluster's replicas should always exist and function
337    /// properly. If it doesn't, however, that is less urgent than
338    /// would be the case for a `SystemCritical` replica.
339    System,
340    /// The cluster is controlled by the user, and might go down for
341    /// reasons outside our control (e.g., OOMs).
342    User,
343}
344
345/// The location of an unmanaged replica.
346#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
347pub struct UnmanagedReplicaLocation {
348    /// The network addresses of the storagectl endpoints for each process in
349    /// the replica.
350    pub storagectl_addrs: Vec<String>,
351    /// The network addresses of the computectl endpoints for each process in
352    /// the replica.
353    pub computectl_addrs: Vec<String>,
354}
355
356/// The location of a managed replica.
357#[derive(Clone, Debug, Serialize, PartialEq)]
358pub struct ManagedReplicaLocation {
359    /// The resource allocation for the replica.
360    pub allocation: ReplicaAllocation,
361    /// SQL size parameter used for allocation
362    pub size: String,
363    /// If `true`, Materialize support owns this replica.
364    pub internal: bool,
365    /// Optional SQL size parameter used for billing.
366    pub billed_as: Option<String>,
367    /// The availability zones the replica may be placed in; empty means
368    /// unconstrained.
369    ///
370    /// For a replica of a managed cluster this is the cluster's
371    /// `AVAILABILITY ZONES` pool; for a replica of an unmanaged cluster it is
372    /// the single user-pinned `AVAILABILITY ZONE`, as a zero- or one-element
373    /// list.
374    ///
375    /// Not serialized: this is re-derived from the cluster config at
376    /// concretization, not read back from a durable record.
377    #[serde(skip)]
378    pub availability_zones: Vec<String>,
379    /// See [`ReplicaLocation::pending`].
380    pub pending: bool,
381}
382
383impl ManagedReplicaLocation {
384    /// Return the size which should be used to determine billing-related information.
385    pub fn size_for_billing(&self) -> &str {
386        self.billed_as.as_deref().unwrap_or(&self.size)
387    }
388}
389
390/// Configures logging for a cluster replica.
391pub type ReplicaLogging = ComputeReplicaLogging;
392
393/// Identifier of a process within a replica.
394pub type ProcessId = u64;
395
396/// An event describing a change in status of a cluster replica process.
397#[derive(Debug, Clone, Serialize)]
398pub struct ClusterEvent {
399    pub cluster_id: ClusterId,
400    pub replica_id: ReplicaId,
401    pub process_id: ProcessId,
402    pub status: ClusterStatus,
403    /// Cumulative restart count of the process, propagated from the orchestrator.
404    /// See [`mz_orchestrator::ServiceEvent::restart_count`].
405    pub restart_count: u64,
406    pub time: DateTime<Utc>,
407}
408
409impl Controller {
410    /// Creates a cluster with the specified identifier and configuration.
411    ///
412    /// A cluster is a combination of a storage instance and a compute instance.
413    /// A cluster has zero or more replicas; each replica colocates the storage
414    /// and compute layers on the same physical resources.
415    pub fn create_cluster(
416        &mut self,
417        id: ClusterId,
418        config: ClusterConfig,
419    ) -> Result<(), anyhow::Error> {
420        self.storage
421            .create_instance(id, config.workload_class.clone());
422        self.compute
423            .create_instance(id, config.arranged_logs, config.workload_class)?;
424        Ok(())
425    }
426
427    /// Updates the workload class for a cluster.
428    ///
429    /// # Panics
430    ///
431    /// Panics if the instance does not exist in the StorageController or the ComputeController.
432    pub fn update_cluster_workload_class(&mut self, id: ClusterId, workload_class: Option<String>) {
433        self.storage
434            .update_instance_workload_class(id, workload_class.clone());
435        self.compute
436            .update_instance_workload_class(id, workload_class)
437            .expect("instance exists");
438    }
439
440    /// Drops the specified cluster.
441    ///
442    /// # Panics
443    ///
444    /// Panics if the cluster still has replicas.
445    pub fn drop_cluster(&mut self, id: ClusterId) {
446        self.storage.drop_instance(id);
447        self.compute.drop_instance(id);
448    }
449
450    /// Creates a replica of the specified cluster with the specified identifier
451    /// and configuration.
452    pub fn create_replica(
453        &mut self,
454        cluster_id: ClusterId,
455        replica_id: ReplicaId,
456        cluster_name: String,
457        replica_name: String,
458        role: ClusterRole,
459        config: ReplicaConfig,
460        enable_worker_core_affinity: bool,
461    ) -> Result<(), anyhow::Error> {
462        let storage_location: ClusterReplicaLocation;
463        let compute_location: ClusterReplicaLocation;
464        let metrics_task: Option<AbortOnDropHandle<()>>;
465
466        match config.location {
467            ReplicaLocation::Unmanaged(UnmanagedReplicaLocation {
468                storagectl_addrs,
469                computectl_addrs,
470            }) => {
471                compute_location = ClusterReplicaLocation {
472                    ctl_addrs: computectl_addrs,
473                };
474                storage_location = ClusterReplicaLocation {
475                    ctl_addrs: storagectl_addrs,
476                };
477                metrics_task = None;
478            }
479            ReplicaLocation::Managed(m) => {
480                let (service, metrics_task_join_handle) = self.provision_replica(
481                    cluster_id,
482                    replica_id,
483                    cluster_name,
484                    replica_name,
485                    role,
486                    m,
487                    enable_worker_core_affinity,
488                )?;
489                storage_location = ClusterReplicaLocation {
490                    ctl_addrs: service.addresses("storagectl"),
491                };
492                compute_location = ClusterReplicaLocation {
493                    ctl_addrs: service.addresses("computectl"),
494                };
495                metrics_task = Some(metrics_task_join_handle);
496
497                // Register the replica for HTTP proxying.
498                let http_addresses = service.addresses("internal-http");
499                self.replica_http_locator
500                    .register_replica(cluster_id, replica_id, http_addresses);
501            }
502        }
503
504        self.storage
505            .connect_replica(cluster_id, replica_id, storage_location);
506        self.compute.add_replica_to_instance(
507            cluster_id,
508            replica_id,
509            compute_location,
510            config.compute,
511        )?;
512
513        if let Some(task) = metrics_task {
514            self.metrics_tasks.insert(replica_id, task);
515        }
516
517        Ok(())
518    }
519
520    /// Drops the specified replica of the specified cluster.
521    pub fn drop_replica(
522        &mut self,
523        cluster_id: ClusterId,
524        replica_id: ReplicaId,
525    ) -> Result<(), anyhow::Error> {
526        // We unconditionally deprovision even for unmanaged replicas to avoid
527        // needing to keep track of which replicas are managed and which are
528        // unmanaged. Deprovisioning is a no-op if the replica ID was never
529        // provisioned.
530        self.deprovision_replica(cluster_id, replica_id, self.deploy_generation)?;
531        self.metrics_tasks.remove(&replica_id);
532
533        // Remove HTTP addresses from the locator.
534        self.replica_http_locator
535            .remove_replica(cluster_id, replica_id);
536
537        // The coordinator only re-pushes the override map when the scoped
538        // configuration itself changes, so a dropped replica's entry would
539        // otherwise be retained until the next such change.
540        self.replica_dyncfg_overrides.remove(&replica_id);
541
542        self.compute.drop_replica(cluster_id, replica_id)?;
543        self.storage.drop_replica(cluster_id, replica_id);
544        Ok(())
545    }
546
547    /// Removes replicas from past generations in a background task.
548    pub(crate) fn remove_past_generation_replicas_in_background(&self) {
549        let deploy_generation = self.deploy_generation;
550        let dyncfg = Arc::clone(self.compute.dyncfg());
551        let orchestrator = Arc::clone(&self.orchestrator);
552        task::spawn(
553            || "controller_remove_past_generation_replicas",
554            async move {
555                info!("attempting to remove past generation replicas");
556                loop {
557                    match try_remove_past_generation_replicas(&*orchestrator, deploy_generation)
558                        .await
559                    {
560                        Ok(()) => {
561                            info!("successfully removed past generation replicas");
562                            return;
563                        }
564                        Err(e) => {
565                            let interval =
566                                CONTROLLER_PAST_GENERATION_REPLICA_CLEANUP_RETRY_INTERVAL
567                                    .get(&dyncfg);
568                            warn!(%e, "failed to remove past generation replicas; will retry in {interval:?}");
569                            time::sleep(interval).await;
570                        }
571                    }
572                }
573            },
574        );
575    }
576
577    /// Remove replicas that are orphaned in the current generation.
578    #[instrument]
579    pub async fn remove_orphaned_replicas(
580        &mut self,
581        next_user_replica_id: u64,
582        next_system_replica_id: u64,
583    ) -> Result<(), anyhow::Error> {
584        let desired: BTreeSet<_> = self.metrics_tasks.keys().copied().collect();
585
586        let actual: BTreeSet<_> = self
587            .orchestrator
588            .list_services()
589            .await?
590            .iter()
591            .map(|s| ReplicaServiceName::from_str(s))
592            .collect::<Result<_, _>>()?;
593
594        for ReplicaServiceName {
595            cluster_id,
596            replica_id,
597            generation,
598        } in actual
599        {
600            // We limit our attention here to replicas from the current deploy
601            // generation. Replicas from past generations are cleaned up during
602            // `Controller::allow_writes`.
603            if generation != self.deploy_generation {
604                continue;
605            }
606
607            let smaller_next = match replica_id {
608                ReplicaId::User(id) if id >= next_user_replica_id => {
609                    Some(ReplicaId::User(next_user_replica_id))
610                }
611                ReplicaId::System(id) if id >= next_system_replica_id => {
612                    Some(ReplicaId::System(next_system_replica_id))
613                }
614                _ => None,
615            };
616            if let Some(next) = smaller_next {
617                // Found a replica in the orchestrator with a higher replica ID
618                // than what we are aware of. This must have been created by an
619                // environmentd that's competing for control of this generation.
620                // Abort to let the other process have full control.
621                halt!("found replica ID ({replica_id}) in orchestrator >= next ID ({next})");
622            }
623            if !desired.contains(&replica_id) {
624                self.deprovision_replica(cluster_id, replica_id, generation)?;
625            }
626        }
627
628        self.orchestrator.flush().await?;
629        Ok(())
630    }
631
632    pub fn events_stream(&self) -> BoxStream<'static, ClusterEvent> {
633        let deploy_generation = self.deploy_generation;
634
635        fn translate_event(event: ServiceEvent) -> Result<(ClusterEvent, u64), anyhow::Error> {
636            let ReplicaServiceName {
637                cluster_id,
638                replica_id,
639                generation: replica_generation,
640                ..
641            } = event.service_id.parse()?;
642
643            let event = ClusterEvent {
644                cluster_id,
645                replica_id,
646                process_id: event.process_id,
647                status: event.status,
648                restart_count: event.restart_count,
649                time: event.time,
650            };
651
652            Ok((event, replica_generation))
653        }
654
655        let stream = self
656            .orchestrator
657            .watch_services()
658            .map(|event| event.and_then(translate_event))
659            .filter_map(move |event| async move {
660                match event {
661                    Ok((event, replica_generation)) => {
662                        if replica_generation == deploy_generation {
663                            Some(event)
664                        } else {
665                            None
666                        }
667                    }
668                    Err(error) => {
669                        error!("service watch error: {error}");
670                        None
671                    }
672                }
673            });
674
675        Box::pin(stream)
676    }
677
678    /// Provisions a replica with the service orchestrator.
679    fn provision_replica(
680        &self,
681        cluster_id: ClusterId,
682        replica_id: ReplicaId,
683        cluster_name: String,
684        replica_name: String,
685        role: ClusterRole,
686        location: ManagedReplicaLocation,
687        enable_worker_core_affinity: bool,
688    ) -> Result<(Box<dyn Service>, AbortOnDropHandle<()>), anyhow::Error> {
689        let service_name = ReplicaServiceName {
690            cluster_id,
691            replica_id,
692            generation: self.deploy_generation,
693        }
694        .to_string();
695        let role_label = match role {
696            ClusterRole::SystemCritical => "system-critical",
697            ClusterRole::System => "system",
698            ClusterRole::User => "user",
699        };
700        let environment_id = self.connection_context().environment_id.clone();
701        let aws_external_id_prefix = self.connection_context().aws_external_id_prefix.clone();
702        let aws_connection_role_arn = self.connection_context().aws_connection_role_arn.clone();
703        let persist_pubsub_url = self.persist_pubsub_url.clone();
704        let secrets_args = self.secrets_args.to_flags();
705
706        // These configure the replica's process rather than environmentd's, so
707        // they are `ParameterScope::Replica` and must be read through this
708        // replica's scoped overrides. They are baked into the process
709        // configuration at provisioning time, so a later change to either the
710        // environment-wide value or the override reaches the replica only when
711        // it is next provisioned.
712        let overrides = self.replica_dyncfg_overrides.get(&replica_id);
713        // Storage and compute arrangements share one maintenance policy, so a
714        // unified replica runs both kinds of arrangement under the same reach.
715        let arrangement_exert_proportionality =
716            ARRANGEMENT_EXERT_PROPORTIONALITY.get_with_overrides(&self.dyncfg, overrides);
717        let storage_proto_timely_config = TimelyConfig {
718            arrangement_exert_proportionality,
719            ..Default::default()
720        };
721        let compute_proto_timely_config = TimelyConfig {
722            arrangement_exert_proportionality,
723            enable_zero_copy: ENABLE_TIMELY_ZERO_COPY.get_with_overrides(&self.dyncfg, overrides),
724            enable_zero_copy_lgalloc: ENABLE_TIMELY_ZERO_COPY_LGALLOC
725                .get_with_overrides(&self.dyncfg, overrides),
726            zero_copy_limit: TIMELY_ZERO_COPY_LIMIT.get_with_overrides(&self.dyncfg, overrides),
727            ..Default::default()
728        };
729        let unified_cluster = ENABLE_UNIFIED_CLUSTER.get_with_overrides(&self.dyncfg, overrides);
730
731        let mut disk_limit = location.allocation.disk_limit;
732        let memory_limit = location.allocation.memory_limit;
733        let mut memory_request = None;
734
735        if location.allocation.swap_enabled {
736            // The disk limit we specify in the service config decides whether or not the replica
737            // gets a scratch disk attached. We want to avoid attaching disks to swap replicas, so
738            // make sure to set the disk limit accordingly.
739            disk_limit = Some(DiskLimit::ZERO);
740
741            // We want to keep the memory request equal to the memory limit, to avoid
742            // over-provisioning and ensure replicas have predictable performance. However, to
743            // enable swap, Kubernetes currently requires that request and limit are different.
744            memory_request = memory_limit.map(|MemoryLimit(limit)| {
745                let request = ByteSize::b(limit.as_u64() - 1);
746                MemoryLimit(request)
747            });
748        }
749
750        let service = self.orchestrator.ensure_service(
751            &service_name,
752            ServiceConfig {
753                app_name: "clusterd".into(),
754                image: self.clusterd_image.clone(),
755                init_container_image: self.init_container_image.clone(),
756                args: Box::new(move |assigned| {
757                    let storage_timely_config = TimelyConfig {
758                        workers: location.allocation.workers.get(),
759                        addresses: assigned.peer_addresses("storage"),
760                        ..storage_proto_timely_config
761                    };
762                    let compute_timely_config = TimelyConfig {
763                        workers: location.allocation.workers.get(),
764                        addresses: assigned.peer_addresses("compute"),
765                        ..compute_proto_timely_config
766                    };
767
768                    let mut args = vec![
769                        format!(
770                            "--storage-controller-listen-addr={}",
771                            assigned.listen_addrs["storagectl"]
772                        ),
773                        format!(
774                            "--compute-controller-listen-addr={}",
775                            assigned.listen_addrs["computectl"]
776                        ),
777                        format!(
778                            "--internal-http-listen-addr={}",
779                            assigned.listen_addrs["internal-http"]
780                        ),
781                        format!("--opentelemetry-resource=cluster_id={}", cluster_id),
782                        format!("--opentelemetry-resource=replica_id={}", replica_id),
783                        format!("--persist-pubsub-url={}", persist_pubsub_url),
784                        format!("--environment-id={}", environment_id),
785                        format!(
786                            "--storage-timely-config={}",
787                            storage_timely_config.to_string(),
788                        ),
789                        format!(
790                            "--compute-timely-config={}",
791                            compute_timely_config.to_string(),
792                        ),
793                    ];
794                    if let Some(aws_external_id_prefix) = &aws_external_id_prefix {
795                        args.push(format!(
796                            "--aws-external-id-prefix={}",
797                            aws_external_id_prefix
798                        ));
799                    }
800                    if let Some(aws_connection_role_arn) = &aws_connection_role_arn {
801                        args.push(format!(
802                            "--aws-connection-role-arn={}",
803                            aws_connection_role_arn
804                        ));
805                    }
806                    if let Some(memory_limit) = location.allocation.memory_limit {
807                        args.push(format!(
808                            "--announce-memory-limit={}",
809                            memory_limit.0.as_u64()
810                        ));
811                    }
812                    if location.allocation.cpu_exclusive && enable_worker_core_affinity {
813                        args.push("--worker-core-affinity".into());
814                    }
815                    if unified_cluster {
816                        args.push("--unified-cluster".into());
817                    }
818                    if location.allocation.is_cc {
819                        args.push("--is-cc".into());
820                    }
821
822                    // If swap is enabled, make the replica limit its own heap usage based on the
823                    // configured memory and disk limits.
824                    if location.allocation.swap_enabled
825                        && let Some(memory_limit) = location.allocation.memory_limit
826                        && let Some(disk_limit) = location.allocation.disk_limit
827                        // Currently, the way for replica sizes to request unlimited swap is to
828                        // specify a `disk_limit` of 0. Ideally we'd change this to make them
829                        // specify no disk limit instead, but for now we need to special-case here.
830                        && disk_limit != DiskLimit::ZERO
831                    {
832                        let heap_limit = memory_limit.0 + disk_limit.0;
833                        args.push(format!("--heap-limit={}", heap_limit.as_u64()));
834                    }
835
836                    args.extend(secrets_args.clone());
837                    args
838                }),
839                ports: vec![
840                    ServicePort {
841                        name: "storagectl".into(),
842                        port_hint: 2100,
843                    },
844                    // To simplify the changes to tests, the port
845                    // chosen here is _after_ the compute ones.
846                    // TODO(petrosagg): fix the numerical ordering here
847                    ServicePort {
848                        name: "storage".into(),
849                        port_hint: 2103,
850                    },
851                    ServicePort {
852                        name: "computectl".into(),
853                        port_hint: 2101,
854                    },
855                    ServicePort {
856                        name: "compute".into(),
857                        port_hint: 2102,
858                    },
859                    ServicePort {
860                        name: "internal-http".into(),
861                        port_hint: 6878,
862                    },
863                ],
864                cpu_limit: location.allocation.cpu_limit,
865                cpu_request: location.allocation.cpu_request,
866                memory_limit,
867                memory_request,
868                scale: location.allocation.scale,
869                labels: BTreeMap::from([
870                    ("replica-id".into(), replica_id.to_string()),
871                    ("cluster-id".into(), cluster_id.to_string()),
872                    ("generation".into(), self.deploy_generation.to_string()),
873                    ("type".into(), "cluster".into()),
874                    ("replica-role".into(), role_label.into()),
875                    ("workers".into(), location.allocation.workers.to_string()),
876                    (
877                        "size".into(),
878                        location
879                            .size
880                            .to_string()
881                            .replace("=", "-")
882                            .replace(",", "_"),
883                    ),
884                ]),
885                annotations: BTreeMap::from([
886                    (
887                        "replica-name".into(),
888                        format!("{cluster_name}.{replica_name}"),
889                    ),
890                    ("cluster-name".into(), cluster_name),
891                ]),
892                // An empty list means no AZ constraint; a non-empty one pins
893                // placement to those zones.
894                availability_zones: Some(location.availability_zones).filter(|azs| !azs.is_empty()),
895                // This provides the orchestrator with some label selectors that
896                // are used to constraint the scheduling of replicas, based on
897                // its internal configuration.
898                //
899                // Selectors include `generation` so that scheduling constraints
900                // (anti-affinity, topology spread) only consider pods of the same
901                // deploy generation. Otherwise, during a generation rollout, the
902                // new-generation pods would be constrained by the placement of
903                // old-generation pods that are about to be torn down, which can
904                // prevent the new pods from scheduling (e.g., when only one AZ
905                // has capacity but it is already occupied by an old-generation
906                // pod).
907                other_replicas_selector: vec![
908                    LabelSelector {
909                        label_name: "cluster-id".to_string(),
910                        logic: LabelSelectionLogic::Eq {
911                            value: cluster_id.to_string(),
912                        },
913                    },
914                    // Select other replicas (but not oneself)
915                    LabelSelector {
916                        label_name: "replica-id".into(),
917                        logic: LabelSelectionLogic::NotEq {
918                            value: replica_id.to_string(),
919                        },
920                    },
921                    LabelSelector {
922                        label_name: "generation".into(),
923                        logic: LabelSelectionLogic::Eq {
924                            value: self.deploy_generation.to_string(),
925                        },
926                    },
927                ],
928                replicas_selector: vec![
929                    LabelSelector {
930                        label_name: "cluster-id".to_string(),
931                        // Select ALL replicas.
932                        logic: LabelSelectionLogic::Eq {
933                            value: cluster_id.to_string(),
934                        },
935                    },
936                    LabelSelector {
937                        label_name: "generation".into(),
938                        logic: LabelSelectionLogic::Eq {
939                            value: self.deploy_generation.to_string(),
940                        },
941                    },
942                ],
943                disk_limit,
944                node_selector: location.allocation.selectors,
945            },
946        )?;
947
948        let metrics_task = mz_ore::task::spawn(|| format!("replica-metrics-{replica_id}"), {
949            let tx = self.metrics_tx.clone();
950            let orchestrator = Arc::clone(&self.orchestrator);
951            let service_name = service_name.clone();
952            async move {
953                const METRICS_INTERVAL: Duration = Duration::from_secs(60);
954
955                // TODO[btv] -- I tried implementing a `watch_metrics` function,
956                // similar to `watch_services`, but it crashed due to
957                // https://github.com/kube-rs/kube/issues/1092 .
958                //
959                // If `metrics-server` can be made to fill in `resourceVersion`,
960                // or if that bug is fixed, we can try that again rather than using this inelegant
961                // loop.
962                let mut interval = tokio::time::interval(METRICS_INTERVAL);
963                loop {
964                    interval.tick().await;
965                    match orchestrator.fetch_service_metrics(&service_name).await {
966                        Ok(metrics) => {
967                            let _ = tx.send((replica_id, metrics));
968                        }
969                        Err(e) => {
970                            warn!("failed to get metrics for replica {replica_id}: {e}");
971                        }
972                    }
973                }
974            }
975        });
976
977        Ok((service, metrics_task.abort_on_drop()))
978    }
979
980    /// Deprovisions a replica with the service orchestrator.
981    fn deprovision_replica(
982        &self,
983        cluster_id: ClusterId,
984        replica_id: ReplicaId,
985        generation: u64,
986    ) -> Result<(), anyhow::Error> {
987        let service_name = ReplicaServiceName {
988            cluster_id,
989            replica_id,
990            generation,
991        }
992        .to_string();
993        self.orchestrator.drop_service(&service_name)
994    }
995}
996
997/// Remove all replicas from past generations.
998async fn try_remove_past_generation_replicas(
999    orchestrator: &dyn NamespacedOrchestrator,
1000    deploy_generation: u64,
1001) -> Result<(), anyhow::Error> {
1002    let services: BTreeSet<_> = orchestrator.list_services().await?.into_iter().collect();
1003
1004    for service in services {
1005        let name: ReplicaServiceName = service.parse()?;
1006        if name.generation < deploy_generation {
1007            info!(
1008                cluster_id = %name.cluster_id,
1009                replica_id = %name.replica_id,
1010                "removing past generation replica",
1011            );
1012            orchestrator.drop_service(&service)?;
1013        }
1014    }
1015
1016    Ok(())
1017}
1018
1019/// Represents the name of a cluster replica service in the orchestrator.
1020#[derive(PartialEq, Eq, PartialOrd, Ord)]
1021pub struct ReplicaServiceName {
1022    pub cluster_id: ClusterId,
1023    pub replica_id: ReplicaId,
1024    pub generation: u64,
1025}
1026
1027impl fmt::Display for ReplicaServiceName {
1028    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
1029        let ReplicaServiceName {
1030            cluster_id,
1031            replica_id,
1032            generation,
1033        } = self;
1034        write!(f, "{cluster_id}-replica-{replica_id}-gen-{generation}")
1035    }
1036}
1037
1038impl FromStr for ReplicaServiceName {
1039    type Err = anyhow::Error;
1040
1041    fn from_str(s: &str) -> Result<Self, Self::Err> {
1042        static SERVICE_NAME_RE: LazyLock<Regex> = LazyLock::new(|| {
1043            Regex::new(r"(?-u)^([us]\d+)-replica-([us]\d+)(?:-gen-(\d+))?$").unwrap()
1044        });
1045
1046        let caps = SERVICE_NAME_RE
1047            .captures(s)
1048            .ok_or_else(|| anyhow!("invalid service name: {s}"))?;
1049
1050        Ok(ReplicaServiceName {
1051            cluster_id: caps.get(1).unwrap().as_str().parse().unwrap(),
1052            replica_id: caps.get(2).unwrap().as_str().parse().unwrap(),
1053            // Old versions of Materialize did not include generations in
1054            // replica service names. Synthesize generation 0 if absent.
1055            // TODO: remove this in the next version of Materialize.
1056            generation: caps.get(3).map_or("0", |m| m.as_str()).parse().unwrap(),
1057        })
1058    }
1059}