Skip to main content

mz_orchestrator_kubernetes/
lib.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::collections::BTreeMap;
11use std::future::Future;
12use std::num::NonZero;
13use std::sync::{Arc, Mutex};
14use std::time::{Duration, Instant};
15use std::{env, fmt};
16
17use anyhow::{Context, anyhow, bail};
18use async_trait::async_trait;
19use chrono::DateTime;
20use clap::ValueEnum;
21use cloud_resource_controller::KubernetesResourceReader;
22use futures::TryFutureExt;
23use futures::stream::{BoxStream, StreamExt};
24use k8s_openapi::DeepMerge;
25use k8s_openapi::api::apps::v1::{StatefulSet, StatefulSetSpec, StatefulSetUpdateStrategy};
26use k8s_openapi::api::core::v1::{
27    Affinity, Capabilities, Container, ContainerPort, EnvVar, EnvVarSource, EphemeralVolumeSource,
28    NodeAffinity, NodeSelector, NodeSelectorRequirement, NodeSelectorTerm, ObjectFieldSelector,
29    ObjectReference, PersistentVolumeClaim, PersistentVolumeClaimSpec,
30    PersistentVolumeClaimTemplate, Pod, PodAffinity, PodAffinityTerm, PodAntiAffinity,
31    PodSecurityContext, PodSpec, PodTemplateSpec, PreferredSchedulingTerm, ResourceRequirements,
32    SeccompProfile, Secret, SecurityContext, Service as K8sService, ServicePort, ServiceSpec,
33    Sysctl, Toleration, TopologySpreadConstraint, Volume, VolumeMount, VolumeResourceRequirements,
34    WeightedPodAffinityTerm,
35};
36use k8s_openapi::apimachinery::pkg::api::resource::Quantity;
37use k8s_openapi::apimachinery::pkg::apis::meta::v1::{
38    LabelSelector, LabelSelectorRequirement, OwnerReference,
39};
40use k8s_openapi::jiff::Timestamp;
41use kube::ResourceExt;
42use kube::api::{Api, DeleteParams, ObjectMeta, PartialObjectMetaExt, Patch, PatchParams};
43use kube::client::Client;
44use kube::error::Error as K8sError;
45use kube::runtime::{WatchStreamExt, watcher};
46use maplit::btreemap;
47use mz_cloud_resources::AwsExternalIdPrefix;
48use mz_cloud_resources::crd::vpc_endpoint::v1::VpcEndpoint;
49use mz_orchestrator::{
50    DiskLimit, LabelSelectionLogic, LabelSelector as MzLabelSelector, NamespacedOrchestrator,
51    OfflineReason, Orchestrator, Service, ServiceAssignments, ServiceConfig, ServiceEvent,
52    ServiceProcessMetrics, ServiceStatus, recommended_k8s_labels, scheduling_config::*,
53};
54use mz_ore::cast::CastInto;
55use mz_ore::retry::Retry;
56use mz_ore::task::AbortOnDropHandle;
57use serde::Deserialize;
58use sha2::{Digest, Sha256};
59use tokio::sync::{mpsc, oneshot};
60use tracing::{error, info, warn};
61
62pub mod cloud_resource_controller;
63pub mod secrets;
64pub mod util;
65
66const FIELD_MANAGER: &str = "environmentd";
67const NODE_FAILURE_THRESHOLD_SECONDS: i64 = 30;
68/// How many attempts of a failed worker command log at warn level before later attempts log at
69/// error level. With the retry backoff starting at 125ms, the first error-level attempt comes
70/// about 4s after the first failure, so a Kubernetes API error that clears on retry, such as a
71/// new service account's RBAC that has not yet propagated, does not reach Sentry.
72const RETRY_WARN_ATTEMPTS: usize = 5;
73
74const POD_TEMPLATE_HASH_ANNOTATION: &str = "environmentd.materialize.cloud/pod-template-hash";
75
76/// Configures a [`KubernetesOrchestrator`].
77#[derive(Debug, Clone)]
78pub struct KubernetesOrchestratorConfig {
79    /// The name of a Kubernetes context to use, if the Kubernetes configuration
80    /// is loaded from the local kubeconfig.
81    pub context: String,
82    /// The name of a non-default Kubernetes scheduler to use, if any.
83    pub scheduler_name: Option<String>,
84    /// The name of a `PriorityClass` to assign to services, if any.
85    pub priority_class_name: Option<String>,
86    /// Annotations to install on every service created by the orchestrator.
87    pub service_annotations: BTreeMap<String, String>,
88    /// Labels to install on every service created by the orchestrator.
89    pub service_labels: BTreeMap<String, String>,
90    /// Node selector to install on every service created by the orchestrator.
91    pub service_node_selector: BTreeMap<String, String>,
92    /// Affinity to install on every service created by the orchestrator.
93    pub service_affinity: Option<String>,
94    /// Tolerations to install on every service created by the orchestrator.
95    pub service_tolerations: Option<String>,
96    /// The service account that each service should run as, if any.
97    pub service_account: Option<String>,
98    /// The image pull policy to set for services created by the orchestrator.
99    pub image_pull_policy: KubernetesImagePullPolicy,
100    /// An AWS external ID prefix to use when making AWS operations on behalf
101    /// of the environment.
102    pub aws_external_id_prefix: Option<AwsExternalIdPrefix>,
103    /// Whether to use code coverage mode or not. Always false for production.
104    pub coverage: bool,
105    /// The Kubernetes StorageClass to use for the ephemeral volume attached to
106    /// services that request disk.
107    ///
108    /// If unspecified, the orchestrator will refuse to create services that
109    /// request disk.
110    pub ephemeral_volume_storage_class: Option<String>,
111    /// The optional fs group for service's pods' `securityContext`.
112    pub service_fs_group: Option<i64>,
113    /// The prefix to prepend to all object names
114    pub name_prefix: Option<String>,
115    /// Whether we should attempt to collect metrics from kubernetes
116    pub collect_pod_metrics: bool,
117    /// Whether to annotate pods for prometheus service discovery.
118    pub enable_prometheus_scrape_annotations: bool,
119}
120
121impl KubernetesOrchestratorConfig {
122    pub fn name_prefix(&self) -> String {
123        self.name_prefix.clone().unwrap_or_default()
124    }
125}
126
127/// Specifies whether Kubernetes should pull Docker images when creating pods.
128#[derive(ValueEnum, Debug, Clone, Copy)]
129pub enum KubernetesImagePullPolicy {
130    /// Always pull the Docker image from the registry.
131    Always,
132    /// Pull the Docker image only if the image is not present.
133    IfNotPresent,
134    /// Never pull the Docker image.
135    Never,
136}
137
138impl fmt::Display for KubernetesImagePullPolicy {
139    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
140        match self {
141            KubernetesImagePullPolicy::Always => f.write_str("Always"),
142            KubernetesImagePullPolicy::IfNotPresent => f.write_str("IfNotPresent"),
143            KubernetesImagePullPolicy::Never => f.write_str("Never"),
144        }
145    }
146}
147
148impl KubernetesImagePullPolicy {
149    pub fn as_kebab_case_str(&self) -> &'static str {
150        match self {
151            Self::Always => "always",
152            Self::IfNotPresent => "if-not-present",
153            Self::Never => "never",
154        }
155    }
156}
157
158/// An orchestrator backed by Kubernetes.
159pub struct KubernetesOrchestrator {
160    client: Client,
161    kubernetes_namespace: String,
162    config: KubernetesOrchestratorConfig,
163    secret_api: Api<Secret>,
164    vpc_endpoint_api: Api<VpcEndpoint>,
165    namespaces: Mutex<BTreeMap<String, Arc<dyn NamespacedOrchestrator>>>,
166    resource_reader: Arc<KubernetesResourceReader>,
167}
168
169impl fmt::Debug for KubernetesOrchestrator {
170    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
171        f.debug_struct("KubernetesOrchestrator").finish()
172    }
173}
174
175impl KubernetesOrchestrator {
176    /// Creates a new Kubernetes orchestrator from the provided configuration.
177    pub async fn new(
178        config: KubernetesOrchestratorConfig,
179    ) -> Result<KubernetesOrchestrator, anyhow::Error> {
180        let (client, kubernetes_namespace) = util::create_client(config.context.clone()).await?;
181        let resource_reader =
182            Arc::new(KubernetesResourceReader::new(config.context.clone()).await?);
183        Ok(KubernetesOrchestrator {
184            client: client.clone(),
185            kubernetes_namespace,
186            config,
187            secret_api: Api::default_namespaced(client.clone()),
188            vpc_endpoint_api: Api::default_namespaced(client),
189            namespaces: Mutex::new(BTreeMap::new()),
190            resource_reader,
191        })
192    }
193}
194
195impl Orchestrator for KubernetesOrchestrator {
196    fn namespace(&self, namespace: &str) -> Arc<dyn NamespacedOrchestrator> {
197        let mut namespaces = self.namespaces.lock().expect("lock poisoned");
198        Arc::clone(namespaces.entry(namespace.into()).or_insert_with(|| {
199            let (command_tx, command_rx) = mpsc::unbounded_channel();
200            let worker = OrchestratorWorker {
201                metrics_api: Api::default_namespaced(self.client.clone()),
202                service_api: Api::default_namespaced(self.client.clone()),
203                stateful_set_api: Api::default_namespaced(self.client.clone()),
204                pod_api: Api::default_namespaced(self.client.clone()),
205                owner_references: vec![],
206                command_rx,
207                name_prefix: self.config.name_prefix.clone().unwrap_or_default(),
208                collect_pod_metrics: self.config.collect_pod_metrics,
209            }
210            .spawn(format!("kubernetes-orchestrator-worker:{namespace}"));
211
212            Arc::new(NamespacedKubernetesOrchestrator {
213                pod_api: Api::default_namespaced(self.client.clone()),
214                kubernetes_namespace: self.kubernetes_namespace.clone(),
215                namespace: namespace.into(),
216                config: self.config.clone(),
217                // TODO(guswynn): make this configurable.
218                scheduling_config: Default::default(),
219                service_infos: std::sync::Mutex::new(BTreeMap::new()),
220                command_tx,
221                _worker: worker,
222            })
223        }))
224    }
225}
226
227#[derive(Clone, Copy)]
228struct ServiceInfo {
229    scale: NonZero<u16>,
230}
231
232struct NamespacedKubernetesOrchestrator {
233    pod_api: Api<Pod>,
234    kubernetes_namespace: String,
235    namespace: String,
236    config: KubernetesOrchestratorConfig,
237    scheduling_config: std::sync::RwLock<ServiceSchedulingConfig>,
238    service_infos: std::sync::Mutex<BTreeMap<String, ServiceInfo>>,
239    command_tx: mpsc::UnboundedSender<WorkerCommand>,
240    _worker: AbortOnDropHandle<()>,
241}
242
243impl fmt::Debug for NamespacedKubernetesOrchestrator {
244    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
245        f.debug_struct("NamespacedKubernetesOrchestrator")
246            .field("kubernetes_namespace", &self.kubernetes_namespace)
247            .field("namespace", &self.namespace)
248            .field("config", &self.config)
249            .finish()
250    }
251}
252
253/// Commands sent from a [`NamespacedKubernetesOrchestrator`] to its
254/// [`OrchestratorWorker`].
255///
256/// Commands for which the caller expects a result include a `result_tx` on which the
257/// [`OrchestratorWorker`] will deliver the result.
258enum WorkerCommand {
259    EnsureService {
260        desc: ServiceDescription,
261    },
262    DropService {
263        name: String,
264    },
265    ListServices {
266        namespace: String,
267        result_tx: oneshot::Sender<Vec<String>>,
268    },
269    Flush {
270        result_tx: oneshot::Sender<()>,
271    },
272    FetchServiceMetrics {
273        name: String,
274        info: ServiceInfo,
275        result_tx: oneshot::Sender<Vec<ServiceProcessMetrics>>,
276    },
277}
278
279/// A description of a service to be created by an [`OrchestratorWorker`].
280#[derive(Debug, Clone)]
281struct ServiceDescription {
282    name: String,
283    scale: NonZero<u16>,
284    service: K8sService,
285    stateful_set: StatefulSet,
286    pod_template_hash: String,
287}
288
289/// A task executing blocking work for a [`NamespacedKubernetesOrchestrator`] in the background.
290///
291/// This type exists to enable making [`NamespacedKubernetesOrchestrator::ensure_service`] and
292/// [`NamespacedKubernetesOrchestrator::drop_service`] non-blocking, allowing invocation of these
293/// methods in latency-sensitive contexts.
294///
295/// Note that, apart from `ensure_service` and `drop_service`, this worker also handles blocking
296/// orchestrator calls that query service state (such as `list_services`). These need to be
297/// sequenced through the worker loop to ensure they linearize as expected. For example, we want to
298/// ensure that a `list_services` result contains exactly those services that were previously
299/// created with `ensure_service` and not yet dropped with `drop_service`.
300struct OrchestratorWorker {
301    metrics_api: Api<PodMetrics>,
302    service_api: Api<K8sService>,
303    stateful_set_api: Api<StatefulSet>,
304    pod_api: Api<Pod>,
305    owner_references: Vec<OwnerReference>,
306    command_rx: mpsc::UnboundedReceiver<WorkerCommand>,
307    name_prefix: String,
308    collect_pod_metrics: bool,
309}
310
311#[derive(Deserialize, Clone, Debug)]
312pub struct PodMetricsContainer {
313    pub name: String,
314    pub usage: PodMetricsContainerUsage,
315}
316
317#[derive(Deserialize, Clone, Debug)]
318pub struct PodMetricsContainerUsage {
319    pub cpu: Quantity,
320    pub memory: Quantity,
321}
322
323#[derive(Deserialize, Clone, Debug)]
324pub struct PodMetrics {
325    pub metadata: ObjectMeta,
326    pub timestamp: String,
327    pub window: String,
328    pub containers: Vec<PodMetricsContainer>,
329}
330
331impl k8s_openapi::Resource for PodMetrics {
332    const GROUP: &'static str = "metrics.k8s.io";
333    const KIND: &'static str = "PodMetrics";
334    const VERSION: &'static str = "v1beta1";
335    const API_VERSION: &'static str = "metrics.k8s.io/v1beta1";
336    const URL_PATH_SEGMENT: &'static str = "pods";
337
338    type Scope = k8s_openapi::NamespaceResourceScope;
339}
340
341impl k8s_openapi::Metadata for PodMetrics {
342    type Ty = ObjectMeta;
343
344    fn metadata(&self) -> &Self::Ty {
345        &self.metadata
346    }
347
348    fn metadata_mut(&mut self) -> &mut Self::Ty {
349        &mut self.metadata
350    }
351}
352
353// Note that these types are very weird. We are `get`-ing a
354// `List` object, and lying about it having an `ObjectMeta`
355// (it deserializes as empty, but we don't need it). The custom
356// metrics API is designed this way, which is very non-standard.
357// A discussion in the `kube` channel in the `tokio` discord
358// confirmed that this layout + using `get_subresource` is the
359// best way to handle this.
360
361#[derive(Deserialize, Clone, Debug)]
362pub struct MetricIdentifier {
363    #[serde(rename = "metricName")]
364    pub name: String,
365    // We skip `selector` for now, as we don't use it
366}
367
368#[derive(Deserialize, Clone, Debug)]
369pub struct MetricValue {
370    #[serde(rename = "describedObject")]
371    pub described_object: ObjectReference,
372    #[serde(flatten)]
373    pub metric_identifier: MetricIdentifier,
374    pub timestamp: String,
375    pub value: Quantity,
376    // We skip `windowSeconds`, as we don't need it
377}
378
379impl NamespacedKubernetesOrchestrator {
380    fn service_name(&self, id: &str) -> String {
381        format!(
382            "{}{}-{id}",
383            self.config.name_prefix.as_deref().unwrap_or(""),
384            self.namespace
385        )
386    }
387
388    /// Return a `watcher::Config` instance that limits results to the namespace
389    /// assigned to this orchestrator.
390    fn watch_pod_params(&self) -> watcher::Config {
391        let ns_selector = format!(
392            "environmentd.materialize.cloud/namespace={}",
393            self.namespace
394        );
395        // This watcher timeout must be shorter than the client read timeout.
396        watcher::Config::default().timeout(59).labels(&ns_selector)
397    }
398
399    /// Convert a higher-level label key to the actual one we
400    /// will give to Kubernetes
401    fn make_label_key(&self, key: &str) -> String {
402        format!("{}.environmentd.materialize.cloud/{}", self.namespace, key)
403    }
404
405    fn label_selector_to_k8s(
406        &self,
407        MzLabelSelector { label_name, logic }: MzLabelSelector,
408    ) -> Result<LabelSelectorRequirement, anyhow::Error> {
409        let (operator, values) = match logic {
410            LabelSelectionLogic::Eq { value } => Ok(("In", vec![value])),
411            LabelSelectionLogic::NotEq { value } => Ok(("NotIn", vec![value])),
412            LabelSelectionLogic::Exists => Ok(("Exists", vec![])),
413            LabelSelectionLogic::NotExists => Ok(("DoesNotExist", vec![])),
414            LabelSelectionLogic::InSet { values } => {
415                if values.is_empty() {
416                    Err(anyhow!(
417                        "Invalid selector logic for {label_name}: empty `in` set"
418                    ))
419                } else {
420                    Ok(("In", values))
421                }
422            }
423            LabelSelectionLogic::NotInSet { values } => {
424                if values.is_empty() {
425                    Err(anyhow!(
426                        "Invalid selector logic for {label_name}: empty `notin` set"
427                    ))
428                } else {
429                    Ok(("NotIn", values))
430                }
431            }
432        }?;
433        let lsr = LabelSelectorRequirement {
434            key: self.make_label_key(&label_name),
435            operator: operator.to_string(),
436            values: Some(values),
437        };
438        Ok(lsr)
439    }
440
441    fn send_command(&self, cmd: WorkerCommand) {
442        self.command_tx.send(cmd).expect("worker task not dropped");
443    }
444}
445
446#[derive(Debug)]
447struct ScaledQuantity {
448    integral_part: u64,
449    exponent: i8,
450    base10: bool,
451}
452
453impl ScaledQuantity {
454    pub fn try_to_integer(&self, scale: i8, base10: bool) -> Option<u64> {
455        if base10 != self.base10 {
456            return None;
457        }
458        let exponent = self.exponent - scale;
459        let mut result = self.integral_part;
460        let base = if self.base10 { 10 } else { 2 };
461        if exponent < 0 {
462            for _ in exponent..0 {
463                result /= base;
464            }
465        } else {
466            for _ in 0..exponent {
467                result = result.checked_mul(base)?;
468            }
469        }
470        Some(result)
471    }
472}
473
474// Parse a k8s `Quantity` object
475// into a numeric value.
476//
477// This is intended to support collecting CPU and Memory data.
478// Thus, there are a few that things Kubernetes attempts to do, that we don't,
479// because I've never observed metrics-server specifically sending them:
480// (1) Handle negative numbers (because it's not useful for that use-case)
481// (2) Handle non-integers (because I have never observed them being actually sent)
482// (3) Handle scientific notation (e.g. 1.23e2)
483fn parse_k8s_quantity(s: &str) -> Result<ScaledQuantity, anyhow::Error> {
484    const DEC_SUFFIXES: &[(&str, i8)] = &[
485        ("n", -9),
486        ("u", -6),
487        ("m", -3),
488        ("", 0),
489        ("k", 3), // yep, intentionally lowercase.
490        ("M", 6),
491        ("G", 9),
492        ("T", 12),
493        ("P", 15),
494        ("E", 18),
495    ];
496    const BIN_SUFFIXES: &[(&str, i8)] = &[
497        ("", 0),
498        ("Ki", 10),
499        ("Mi", 20),
500        ("Gi", 30),
501        ("Ti", 40),
502        ("Pi", 50),
503        ("Ei", 60),
504    ];
505
506    let (positive, s) = match s.chars().next() {
507        Some('+') => (true, &s[1..]),
508        Some('-') => (false, &s[1..]),
509        _ => (true, s),
510    };
511
512    if !positive {
513        anyhow::bail!("Negative numbers not supported")
514    }
515
516    fn is_suffix_char(ch: char) -> bool {
517        "numkMGTPEKi".contains(ch)
518    }
519    let (num, suffix) = match s.find(is_suffix_char) {
520        None => (s, ""),
521        Some(idx) => s.split_at(idx),
522    };
523    let num: u64 = num.parse()?;
524    let (exponent, base10) = if let Some((_, exponent)) =
525        DEC_SUFFIXES.iter().find(|(target, _)| suffix == *target)
526    {
527        (exponent, true)
528    } else if let Some((_, exponent)) = BIN_SUFFIXES.iter().find(|(target, _)| suffix == *target) {
529        (exponent, false)
530    } else {
531        anyhow::bail!("Unrecognized suffix: {suffix}");
532    };
533    Ok(ScaledQuantity {
534        integral_part: num,
535        exponent: *exponent,
536        base10,
537    })
538}
539
540#[async_trait]
541impl NamespacedOrchestrator for NamespacedKubernetesOrchestrator {
542    async fn fetch_service_metrics(
543        &self,
544        id: &str,
545    ) -> Result<Vec<ServiceProcessMetrics>, anyhow::Error> {
546        let info = if let Some(info) = self.service_infos.lock().expect("poisoned lock").get(id) {
547            *info
548        } else {
549            // This should have been set in `ensure_service`.
550            tracing::error!("Failed to get info for {id}");
551            anyhow::bail!("Failed to get info for {id}");
552        };
553
554        let (result_tx, result_rx) = oneshot::channel();
555        self.send_command(WorkerCommand::FetchServiceMetrics {
556            name: self.service_name(id),
557            info,
558            result_tx,
559        });
560
561        let metrics = result_rx.await.expect("worker task not dropped");
562        Ok(metrics)
563    }
564
565    fn ensure_service(
566        &self,
567        id: &str,
568        ServiceConfig {
569            app_name,
570            image,
571            init_container_image,
572            args,
573            ports: ports_in,
574            memory_limit,
575            memory_request,
576            cpu_limit,
577            cpu_request,
578            scale,
579            labels: labels_in,
580            annotations: annotations_in,
581            availability_zones,
582            other_replicas_selector,
583            replicas_selector,
584            disk_limit,
585            node_selector,
586        }: ServiceConfig,
587    ) -> Result<Box<dyn Service>, anyhow::Error> {
588        // This is extremely cheap to clone, so just look into the lock once.
589        let scheduling_config: ServiceSchedulingConfig =
590            self.scheduling_config.read().expect("poisoned").clone();
591
592        // Enable disk if the size does not disable it.
593        let disk = disk_limit != Some(DiskLimit::ZERO);
594
595        let name = self.service_name(id);
596        // The match labels should be the minimal set of labels that uniquely
597        // identify the pods in the stateful set. Changing these after the
598        // `StatefulSet` is created is not permitted by Kubernetes, and we're
599        // not yet smart enough to handle deleting and recreating the
600        // `StatefulSet`.
601        let mut match_labels = btreemap! {
602            "environmentd.materialize.cloud/namespace".into() => self.namespace.clone(),
603            "environmentd.materialize.cloud/service-id".into() => id.into(),
604        };
605        for (key, value) in &self.config.service_labels {
606            match_labels.insert(key.clone(), value.clone());
607        }
608
609        let mut labels = match_labels.clone();
610        for (key, value) in labels_in {
611            labels.insert(self.make_label_key(&key), value);
612        }
613
614        let standard_labels = recommended_k8s_labels(app_name);
615
616        // Standard Kubernetes labels
617        labels.extend(standard_labels.clone());
618
619        labels.insert(self.make_label_key("scale"), scale.to_string());
620
621        for port in &ports_in {
622            labels.insert(
623                format!("environmentd.materialize.cloud/port-{}", port.name),
624                "true".into(),
625            );
626        }
627        let mut limits = BTreeMap::new();
628        let mut requests = BTreeMap::new();
629        if let Some(memory_limit) = memory_limit {
630            limits.insert(
631                "memory".into(),
632                Quantity(memory_limit.0.as_u64().to_string()),
633            );
634            requests.insert(
635                "memory".into(),
636                Quantity(memory_limit.0.as_u64().to_string()),
637            );
638        }
639        if let Some(memory_request) = memory_request {
640            requests.insert(
641                "memory".into(),
642                Quantity(memory_request.0.as_u64().to_string()),
643            );
644        }
645        if let Some(cpu_limit) = cpu_limit {
646            limits.insert(
647                "cpu".into(),
648                Quantity(format!("{}m", cpu_limit.as_millicpus())),
649            );
650            requests.insert(
651                "cpu".into(),
652                Quantity(format!("{}m", cpu_limit.as_millicpus())),
653            );
654        }
655        if let Some(cpu_request) = cpu_request {
656            requests.insert(
657                "cpu".into(),
658                Quantity(format!("{}m", cpu_request.as_millicpus())),
659            );
660        }
661        let service = K8sService {
662            metadata: ObjectMeta {
663                name: Some(name.clone()),
664                labels: Some(standard_labels.clone()),
665                ..Default::default()
666            },
667            spec: Some(ServiceSpec {
668                ports: Some(
669                    ports_in
670                        .iter()
671                        .map(|port| ServicePort {
672                            port: port.port_hint.into(),
673                            name: Some(port.name.clone()),
674                            ..Default::default()
675                        })
676                        .collect(),
677                ),
678                cluster_ip: Some("None".to_string()),
679                selector: Some(match_labels.clone()),
680                ..Default::default()
681            }),
682            status: None,
683        };
684
685        let hosts = (0..scale.get())
686            .map(|i| {
687                format!(
688                    "{name}-{i}.{name}.{}.svc.cluster.local",
689                    self.kubernetes_namespace
690                )
691            })
692            .collect::<Vec<_>>();
693        let ports = ports_in
694            .iter()
695            .map(|p| (p.name.clone(), p.port_hint))
696            .collect::<BTreeMap<_, _>>();
697
698        let mut listen_addrs = BTreeMap::new();
699        let mut peer_addrs = vec![BTreeMap::new(); hosts.len()];
700        for (name, port) in &ports {
701            listen_addrs.insert(name.clone(), format!("0.0.0.0:{port}"));
702            for (i, host) in hosts.iter().enumerate() {
703                peer_addrs[i].insert(name.clone(), format!("{host}:{port}"));
704            }
705        }
706        let mut args = args(ServiceAssignments {
707            listen_addrs: &listen_addrs,
708            peer_addrs: &peer_addrs,
709        });
710
711        // This constrains the orchestrator (for those orchestrators that support
712        // anti-affinity, today just k8s) to never schedule pods for different replicas
713        // of the same cluster on the same node. Pods from the _same_ replica are fine;
714        // pods from different clusters are also fine.
715        //
716        // The point is that if pods of two replicas are on the same node, that node
717        // going down would kill both replicas, and so the replication factor of the
718        // cluster in question is illusory.
719        let anti_affinity = Some({
720            let label_selector_requirements = other_replicas_selector
721                .clone()
722                .into_iter()
723                .map(|ls| self.label_selector_to_k8s(ls))
724                .collect::<Result<Vec<_>, _>>()?;
725            let ls = LabelSelector {
726                match_expressions: Some(label_selector_requirements),
727                ..Default::default()
728            };
729            let pat = PodAffinityTerm {
730                label_selector: Some(ls),
731                topology_key: "kubernetes.io/hostname".to_string(),
732                ..Default::default()
733            };
734
735            if !scheduling_config.soften_replication_anti_affinity {
736                PodAntiAffinity {
737                    required_during_scheduling_ignored_during_execution: Some(vec![pat]),
738                    ..Default::default()
739                }
740            } else {
741                PodAntiAffinity {
742                    preferred_during_scheduling_ignored_during_execution: Some(vec![
743                        WeightedPodAffinityTerm {
744                            weight: scheduling_config.soften_replication_anti_affinity_weight,
745                            pod_affinity_term: pat,
746                        },
747                    ]),
748                    ..Default::default()
749                }
750            }
751        });
752
753        let pod_affinity = if let Some(weight) = scheduling_config.multi_pod_az_affinity_weight {
754            // `match_labels` sufficiently selects pods in the same replica.
755            let ls = LabelSelector {
756                match_labels: Some(match_labels.clone()),
757                ..Default::default()
758            };
759            let pat = PodAffinityTerm {
760                label_selector: Some(ls),
761                topology_key: "topology.kubernetes.io/zone".to_string(),
762                ..Default::default()
763            };
764
765            Some(PodAffinity {
766                preferred_during_scheduling_ignored_during_execution: Some(vec![
767                    WeightedPodAffinityTerm {
768                        weight,
769                        pod_affinity_term: pat,
770                    },
771                ]),
772                ..Default::default()
773            })
774        } else {
775            None
776        };
777
778        let topology_spread = if scheduling_config.topology_spread.enabled {
779            let config = &scheduling_config.topology_spread;
780
781            if !config.ignore_non_singular_scale || scale.get() == 1 {
782                let label_selector_requirements = (if config.ignore_non_singular_scale {
783                    let mut replicas_selector_ignoring_scale = replicas_selector.clone();
784
785                    replicas_selector_ignoring_scale.push(mz_orchestrator::LabelSelector {
786                        label_name: "scale".into(),
787                        logic: mz_orchestrator::LabelSelectionLogic::Eq {
788                            value: "1".to_string(),
789                        },
790                    });
791
792                    replicas_selector_ignoring_scale
793                } else {
794                    replicas_selector
795                })
796                .into_iter()
797                .map(|ls| self.label_selector_to_k8s(ls))
798                .collect::<Result<Vec<_>, _>>()?;
799                let ls = LabelSelector {
800                    match_expressions: Some(label_selector_requirements),
801                    ..Default::default()
802                };
803
804                if config.soft && config.min_domains.is_some() {
805                    warn!(
806                        "topology spread is soft but min_domains is set; \
807                         Kubernetes rejects minDomains with ScheduleAnyway, \
808                         so min_domains will be ignored"
809                    );
810                }
811                if availability_zones.is_some() && config.min_domains.is_some() {
812                    warn!(
813                        "topology spread has min_domains set but availability_zones \
814                         constrains eligible topology domains via node affinity; \
815                         minDomains will be ignored to avoid preventing pod scheduling"
816                    );
817                }
818
819                let constraint = TopologySpreadConstraint {
820                    label_selector: Some(ls),
821                    min_domains: topology_spread_min_domains(
822                        config.soft,
823                        availability_zones.is_some(),
824                        config.min_domains,
825                    ),
826                    max_skew: config.max_skew,
827                    topology_key: "topology.kubernetes.io/zone".to_string(),
828                    when_unsatisfiable: if config.soft {
829                        "ScheduleAnyway".to_string()
830                    } else {
831                        "DoNotSchedule".to_string()
832                    },
833                    // TODO(guswynn): restore these once they are supported.
834                    // Consider node affinities when calculating topology spread. This is the
835                    // default: <https://docs.rs/k8s-openapi/latest/k8s_openapi/api/core/v1/struct.TopologySpreadConstraint.html#structfield.node_affinity_policy>,
836                    // made explicit.
837                    // node_affinity_policy: Some("Honor".to_string()),
838                    // Do not consider node taints when calculating topology spread. This is the
839                    // default: <https://docs.rs/k8s-openapi/latest/k8s_openapi/api/core/v1/struct.TopologySpreadConstraint.html#structfield.node_taints_policy>,
840                    // made explicit.
841                    // node_taints_policy: Some("Ignore".to_string()),
842                    match_label_keys: None,
843                    // Once the above are restorted, we should't have `..Default::default()` here because the specifics of these fields are
844                    // subtle enough where we want compilation failures when we upgrade
845                    ..Default::default()
846                };
847                Some(vec![constraint])
848            } else {
849                None
850            }
851        } else {
852            None
853        };
854
855        let mut pod_annotations = btreemap! {
856            // Prevent the cluster-autoscaler (or karpenter) from evicting these pods in attempts to scale down
857            // and terminate nodes.
858            // This will cost us more money, but should give us better uptime.
859            // This does not prevent all evictions by Kubernetes, only the ones initiated by the
860            // cluster-autoscaler (or karpenter). Notably, eviction of pods for resource overuse is still enabled.
861            "cluster-autoscaler.kubernetes.io/safe-to-evict".to_owned() => "false".to_string(),
862            "karpenter.sh/do-not-evict".to_owned() => "true".to_string(),
863
864            // It's called do-not-disrupt in newer versions of karpenter, so adding for forward/backward compatibility
865            "karpenter.sh/do-not-disrupt".to_owned() => "true".to_string(),
866        };
867        for (key, value) in annotations_in {
868            // We want to use the same prefix as our labels keys
869            pod_annotations.insert(self.make_label_key(&key), value);
870        }
871        if self.config.enable_prometheus_scrape_annotations {
872            if let Some(internal_http_port) = ports_in
873                .iter()
874                .find(|port| port.name == "internal-http")
875                .map(|port| port.port_hint.to_string())
876            {
877                // Enable prometheus scrape discovery
878                pod_annotations.insert("prometheus.io/scrape".to_owned(), "true".to_string());
879                pod_annotations.insert("prometheus.io/port".to_owned(), internal_http_port);
880                pod_annotations.insert("prometheus.io/path".to_owned(), "/metrics".to_string());
881                pod_annotations.insert("prometheus.io/scheme".to_owned(), "http".to_string());
882            }
883        }
884        for (key, value) in &self.config.service_annotations {
885            pod_annotations.insert(key.clone(), value.clone());
886        }
887
888        let default_node_selector = if disk {
889            vec![("materialize.cloud/disk".to_string(), disk.to_string())]
890        } else {
891            // if the cluster doesn't require disk, we can omit the selector
892            // allowing it to be scheduled onto nodes with and without the
893            // selector
894            vec![]
895        };
896
897        let node_selector: BTreeMap<String, String> = default_node_selector
898            .into_iter()
899            .chain(self.config.service_node_selector.clone())
900            .chain(node_selector)
901            .collect();
902
903        let node_affinity = if let Some(availability_zones) = availability_zones {
904            let selector = NodeSelectorTerm {
905                match_expressions: Some(vec![NodeSelectorRequirement {
906                    key: "materialize.cloud/availability-zone".to_string(),
907                    operator: "In".to_string(),
908                    values: Some(availability_zones),
909                }]),
910                match_fields: None,
911            };
912
913            if scheduling_config.soften_az_affinity {
914                Some(NodeAffinity {
915                    preferred_during_scheduling_ignored_during_execution: Some(vec![
916                        PreferredSchedulingTerm {
917                            preference: selector,
918                            weight: scheduling_config.soften_az_affinity_weight,
919                        },
920                    ]),
921                    required_during_scheduling_ignored_during_execution: None,
922                })
923            } else {
924                Some(NodeAffinity {
925                    preferred_during_scheduling_ignored_during_execution: None,
926                    required_during_scheduling_ignored_during_execution: Some(NodeSelector {
927                        node_selector_terms: vec![selector],
928                    }),
929                })
930            }
931        } else {
932            None
933        };
934
935        let mut affinity = Affinity {
936            pod_anti_affinity: anti_affinity,
937            pod_affinity,
938            node_affinity,
939            ..Default::default()
940        };
941        if let Some(service_affinity) = &self.config.service_affinity {
942            affinity.merge_from(serde_json::from_str(service_affinity)?);
943        }
944
945        let container_name = image
946            .rsplit_once('/')
947            .and_then(|(_, name_version)| name_version.rsplit_once(':'))
948            .context("`image` is not ORG/NAME:VERSION")?
949            .0
950            .to_string();
951
952        let container_security_context = if scheduling_config.security_context_enabled {
953            Some(SecurityContext {
954                privileged: Some(false),
955                run_as_non_root: Some(true),
956                allow_privilege_escalation: Some(false),
957                seccomp_profile: Some(SeccompProfile {
958                    type_: "RuntimeDefault".to_string(),
959                    ..Default::default()
960                }),
961                capabilities: Some(Capabilities {
962                    drop: Some(vec!["ALL".to_string()]),
963                    ..Default::default()
964                }),
965                ..Default::default()
966            })
967        } else {
968            None
969        };
970
971        let init_containers = init_container_image.map(|image| {
972            vec![Container {
973                name: "init".to_string(),
974                image: Some(image),
975                image_pull_policy: Some(self.config.image_pull_policy.to_string()),
976                resources: Some(ResourceRequirements {
977                    claims: None,
978                    limits: Some(limits.clone()),
979                    requests: Some(requests.clone()),
980                }),
981                security_context: container_security_context.clone(),
982                env: Some(vec![
983                    EnvVar {
984                        name: "MZ_NAMESPACE".to_string(),
985                        value_from: Some(EnvVarSource {
986                            field_ref: Some(ObjectFieldSelector {
987                                field_path: "metadata.namespace".to_string(),
988                                ..Default::default()
989                            }),
990                            ..Default::default()
991                        }),
992                        ..Default::default()
993                    },
994                    EnvVar {
995                        name: "MZ_POD_NAME".to_string(),
996                        value_from: Some(EnvVarSource {
997                            field_ref: Some(ObjectFieldSelector {
998                                field_path: "metadata.name".to_string(),
999                                ..Default::default()
1000                            }),
1001                            ..Default::default()
1002                        }),
1003                        ..Default::default()
1004                    },
1005                    EnvVar {
1006                        name: "MZ_NODE_NAME".to_string(),
1007                        value_from: Some(EnvVarSource {
1008                            field_ref: Some(ObjectFieldSelector {
1009                                field_path: "spec.nodeName".to_string(),
1010                                ..Default::default()
1011                            }),
1012                            ..Default::default()
1013                        }),
1014                        ..Default::default()
1015                    },
1016                ]),
1017                ..Default::default()
1018            }]
1019        });
1020
1021        let env = if self.config.coverage {
1022            Some(vec![EnvVar {
1023                name: "LLVM_PROFILE_FILE".to_string(),
1024                value: Some(format!("/coverage/{}-%p-%9m%c.profraw", self.namespace)),
1025                ..Default::default()
1026            }])
1027        } else {
1028            None
1029        };
1030
1031        let mut volume_mounts = vec![];
1032
1033        if self.config.coverage {
1034            volume_mounts.push(VolumeMount {
1035                name: "coverage".to_string(),
1036                mount_path: "/coverage".to_string(),
1037                ..Default::default()
1038            })
1039        }
1040
1041        let volumes = match (disk, &self.config.ephemeral_volume_storage_class) {
1042            (true, Some(ephemeral_volume_storage_class)) => {
1043                volume_mounts.push(VolumeMount {
1044                    name: "scratch".to_string(),
1045                    mount_path: "/scratch".to_string(),
1046                    ..Default::default()
1047                });
1048                args.push("--scratch-directory=/scratch".into());
1049
1050                Some(vec![Volume {
1051                    name: "scratch".to_string(),
1052                    ephemeral: Some(EphemeralVolumeSource {
1053                        volume_claim_template: Some(PersistentVolumeClaimTemplate {
1054                            spec: PersistentVolumeClaimSpec {
1055                                access_modes: Some(vec!["ReadWriteOnce".to_string()]),
1056                                storage_class_name: Some(
1057                                    ephemeral_volume_storage_class.to_string(),
1058                                ),
1059                                resources: Some(VolumeResourceRequirements {
1060                                    requests: Some(BTreeMap::from([(
1061                                        "storage".to_string(),
1062                                        Quantity(
1063                                            disk_limit
1064                                                .unwrap_or(DiskLimit::ARBITRARY)
1065                                                .0
1066                                                .as_u64()
1067                                                .to_string(),
1068                                        ),
1069                                    )])),
1070                                    ..Default::default()
1071                                }),
1072                                ..Default::default()
1073                            },
1074                            ..Default::default()
1075                        }),
1076                        ..Default::default()
1077                    }),
1078                    ..Default::default()
1079                }])
1080            }
1081            (true, None) => {
1082                return Err(anyhow!(
1083                    "service requested disk but no ephemeral volume storage class was configured"
1084                ));
1085            }
1086            (false, _) => None,
1087        };
1088
1089        if let Some(name_prefix) = &self.config.name_prefix {
1090            args.push(format!("--secrets-reader-name-prefix={}", name_prefix));
1091        }
1092
1093        let volume_claim_templates = if self.config.coverage {
1094            Some(vec![PersistentVolumeClaim {
1095                metadata: ObjectMeta {
1096                    name: Some("coverage".to_string()),
1097                    ..Default::default()
1098                },
1099                spec: Some(PersistentVolumeClaimSpec {
1100                    access_modes: Some(vec!["ReadWriteOnce".to_string()]),
1101                    resources: Some(VolumeResourceRequirements {
1102                        requests: Some(BTreeMap::from([(
1103                            "storage".to_string(),
1104                            Quantity("10Gi".to_string()),
1105                        )])),
1106                        ..Default::default()
1107                    }),
1108                    ..Default::default()
1109                }),
1110                ..Default::default()
1111            }])
1112        } else {
1113            None
1114        };
1115
1116        let tcp_keepalive_sysctls = vec![
1117            Sysctl {
1118                name: "net.ipv4.tcp_keepalive_time".to_string(),
1119                value: "300".to_string(),
1120            },
1121            Sysctl {
1122                name: "net.ipv4.tcp_keepalive_intvl".to_string(),
1123                value: "30".to_string(),
1124            },
1125            Sysctl {
1126                name: "net.ipv4.tcp_keepalive_probes".to_string(),
1127                value: "3".to_string(),
1128            },
1129        ];
1130
1131        let security_context = if let Some(fs_group) = self.config.service_fs_group {
1132            Some(PodSecurityContext {
1133                fs_group: Some(fs_group),
1134                run_as_user: Some(fs_group),
1135                run_as_group: Some(fs_group),
1136                sysctls: Some(tcp_keepalive_sysctls),
1137                ..Default::default()
1138            })
1139        } else {
1140            Some(PodSecurityContext {
1141                sysctls: Some(tcp_keepalive_sysctls),
1142                ..Default::default()
1143            })
1144        };
1145
1146        let mut tolerations = vec![
1147            // When the node becomes `NotReady` it indicates there is a problem
1148            // with the node. By default Kubernetes waits 300s (5 minutes)
1149            // before descheduling the pod, but we tune this to 30s for faster
1150            // recovery in the case of node failure.
1151            Toleration {
1152                effect: Some("NoExecute".into()),
1153                key: Some("node.kubernetes.io/not-ready".into()),
1154                operator: Some("Exists".into()),
1155                toleration_seconds: Some(NODE_FAILURE_THRESHOLD_SECONDS),
1156                value: None,
1157            },
1158            Toleration {
1159                effect: Some("NoExecute".into()),
1160                key: Some("node.kubernetes.io/unreachable".into()),
1161                operator: Some("Exists".into()),
1162                toleration_seconds: Some(NODE_FAILURE_THRESHOLD_SECONDS),
1163                value: None,
1164            },
1165        ];
1166        if let Some(service_tolerations) = &self.config.service_tolerations {
1167            tolerations.extend(serde_json::from_str::<Vec<_>>(service_tolerations)?);
1168        }
1169        let tolerations = Some(tolerations);
1170
1171        let mut pod_template_spec = PodTemplateSpec {
1172            metadata: Some(ObjectMeta {
1173                labels: Some(labels.clone()),
1174                // Only set `annotations` _after_ we have computed the pod template hash, to
1175                // avoid that annotation changes cause pod replacements.
1176                ..Default::default()
1177            }),
1178            spec: Some(PodSpec {
1179                init_containers,
1180                containers: vec![Container {
1181                    name: container_name,
1182                    image: Some(image),
1183                    args: Some(args),
1184                    image_pull_policy: Some(self.config.image_pull_policy.to_string()),
1185                    ports: Some(
1186                        ports_in
1187                            .iter()
1188                            .map(|port| ContainerPort {
1189                                container_port: port.port_hint.into(),
1190                                name: Some(port.name.clone()),
1191                                ..Default::default()
1192                            })
1193                            .collect(),
1194                    ),
1195                    security_context: container_security_context.clone(),
1196                    resources: Some(ResourceRequirements {
1197                        claims: None,
1198                        limits: Some(limits),
1199                        requests: Some(requests),
1200                    }),
1201                    volume_mounts: if !volume_mounts.is_empty() {
1202                        Some(volume_mounts)
1203                    } else {
1204                        None
1205                    },
1206                    env,
1207                    ..Default::default()
1208                }],
1209                volumes,
1210                security_context,
1211                node_selector: Some(node_selector),
1212                scheduler_name: self.config.scheduler_name.clone(),
1213                priority_class_name: self.config.priority_class_name.clone(),
1214                service_account: self.config.service_account.clone(),
1215                affinity: Some(affinity),
1216                topology_spread_constraints: topology_spread,
1217                tolerations,
1218                // Setting a 0s termination grace period has the side effect of
1219                // automatically starting a new pod when the previous pod is
1220                // currently terminating. This enables recovery from a node
1221                // failure with no manual intervention. Without this setting,
1222                // the StatefulSet controller will refuse to start a new pod
1223                // until the failed node is manually removed from the Kubernetes
1224                // cluster.
1225                //
1226                // The Kubernetes documentation strongly advises against this
1227                // setting, as StatefulSets attempt to provide "at most once"
1228                // semantics [0]--that is, the guarantee that for a given pod in
1229                // a StatefulSet there is *at most* one pod with that identity
1230                // running in the cluster.
1231                //
1232                // Materialize services, however, are carefully designed to
1233                // *not* rely on this guarantee. In fact, we do not believe that
1234                // correct distributed systems can meaningfully rely on
1235                // Kubernetes's guarantee--network packets from a pod can be
1236                // arbitrarily delayed, long past that pod's termination.
1237                //
1238                // [0]: https://kubernetes.io/docs/tasks/run-application/force-delete-stateful-set-pod/#statefulset-considerations
1239                termination_grace_period_seconds: Some(0),
1240                ..Default::default()
1241            }),
1242        };
1243        let pod_template_json = serde_json::to_string(&pod_template_spec).unwrap();
1244        let mut hasher = Sha256::new();
1245        hasher.update(pod_template_json);
1246        let pod_template_hash = hex::encode(hasher.finalize());
1247        pod_annotations.insert(
1248            POD_TEMPLATE_HASH_ANNOTATION.to_owned(),
1249            pod_template_hash.clone(),
1250        );
1251
1252        pod_template_spec.metadata.as_mut().unwrap().annotations = Some(pod_annotations);
1253
1254        let stateful_set = StatefulSet {
1255            metadata: ObjectMeta {
1256                name: Some(name.clone()),
1257                labels: Some(standard_labels.clone()),
1258                ..Default::default()
1259            },
1260            spec: Some(StatefulSetSpec {
1261                selector: LabelSelector {
1262                    match_labels: Some(match_labels),
1263                    ..Default::default()
1264                },
1265                service_name: Some(name.clone()),
1266                replicas: Some(scale.cast_into()),
1267                template: pod_template_spec,
1268                update_strategy: Some(StatefulSetUpdateStrategy {
1269                    type_: Some("OnDelete".to_owned()),
1270                    ..Default::default()
1271                }),
1272                pod_management_policy: Some("Parallel".to_string()),
1273                volume_claim_templates,
1274                ..Default::default()
1275            }),
1276            status: None,
1277        };
1278
1279        self.send_command(WorkerCommand::EnsureService {
1280            desc: ServiceDescription {
1281                name,
1282                scale,
1283                service,
1284                stateful_set,
1285                pod_template_hash,
1286            },
1287        });
1288
1289        self.service_infos
1290            .lock()
1291            .expect("poisoned lock")
1292            .insert(id.to_string(), ServiceInfo { scale });
1293
1294        Ok(Box::new(KubernetesService { hosts, ports }))
1295    }
1296
1297    /// Drops the identified service, if it exists.
1298    fn drop_service(&self, id: &str) -> Result<(), anyhow::Error> {
1299        fail::fail_point!("kubernetes_drop_service", |_| Err(anyhow!("failpoint")));
1300        self.service_infos.lock().expect("poisoned lock").remove(id);
1301
1302        self.send_command(WorkerCommand::DropService {
1303            name: self.service_name(id),
1304        });
1305
1306        Ok(())
1307    }
1308
1309    /// Lists the identifiers of all known services.
1310    async fn list_services(&self) -> Result<Vec<String>, anyhow::Error> {
1311        let (result_tx, result_rx) = oneshot::channel();
1312        self.send_command(WorkerCommand::ListServices {
1313            namespace: self.namespace.clone(),
1314            result_tx,
1315        });
1316
1317        let list = result_rx.await.expect("worker task not dropped");
1318        Ok(list)
1319    }
1320
1321    async fn flush(&self) -> Result<(), anyhow::Error> {
1322        let (result_tx, result_rx) = oneshot::channel();
1323        self.send_command(WorkerCommand::Flush { result_tx });
1324        result_rx.await.expect("worker task not dropped");
1325        Ok(())
1326    }
1327
1328    fn watch_services(&self) -> BoxStream<'static, Result<ServiceEvent, anyhow::Error>> {
1329        fn into_service_event(pod: Pod) -> Result<ServiceEvent, anyhow::Error> {
1330            let process_id = pod.name_any().split('-').next_back().unwrap().parse()?;
1331            let service_id_label = "environmentd.materialize.cloud/service-id";
1332            let service_id = pod
1333                .labels()
1334                .get(service_id_label)
1335                .ok_or_else(|| anyhow!("missing label: {service_id_label}"))?
1336                .clone();
1337
1338            let oomed = pod
1339                .status
1340                .as_ref()
1341                .and_then(|status| status.container_statuses.as_ref())
1342                .map(|container_statuses| {
1343                    container_statuses.iter().any(|cs| {
1344                        // The container might have already transitioned from "terminated" to
1345                        // "waiting"/"running" state, in which case we need to check its previous
1346                        // state to find out why it terminated.
1347                        let current_state = cs.state.as_ref().and_then(|s| s.terminated.as_ref());
1348                        let last_state = cs.last_state.as_ref().and_then(|s| s.terminated.as_ref());
1349                        let termination_state = current_state.or(last_state);
1350
1351                        // The interesting exit codes are:
1352                        //  * 135 (SIGBUS): occurs when lgalloc runs out of disk
1353                        //  * 137 (SIGKILL): occurs when the OOM killer terminates the container
1354                        //  * 167: occurs when the lgalloc or memory limiter terminates the process
1355                        // We treat the all of these as OOM conditions since swap and lgalloc use
1356                        // disk only for spilling memory.
1357                        let exit_code = termination_state.map(|s| s.exit_code);
1358                        exit_code.is_some_and(|e| [135, 137, 167].contains(&e))
1359                    })
1360                })
1361                .unwrap_or(false);
1362
1363            // Sum the per-container restart counts to get a per-process (per-pod)
1364            // restart count. This is cumulative and survives gaps in the watch
1365            // stream, so it lets consumers detect restarts they'd otherwise miss
1366            // by only sampling the ready/not-ready status.
1367            let restart_count = pod
1368                .status
1369                .as_ref()
1370                .and_then(|status| status.container_statuses.as_ref())
1371                .map(|container_statuses| {
1372                    container_statuses
1373                        .iter()
1374                        .map(|cs| u64::try_from(cs.restart_count).unwrap_or(0))
1375                        .sum()
1376                })
1377                .unwrap_or(0);
1378
1379            let (pod_ready, last_probe_time) = pod
1380                .status
1381                .and_then(|status| status.conditions)
1382                .and_then(|conditions| conditions.into_iter().find(|c| c.type_ == "Ready"))
1383                .map(|c| (c.status == "True", c.last_probe_time))
1384                .unwrap_or((false, None));
1385
1386            let status = if pod_ready {
1387                ServiceStatus::Online
1388            } else {
1389                ServiceStatus::Offline(oomed.then_some(OfflineReason::OomKilled))
1390            };
1391            let time = if let Some(time) = last_probe_time {
1392                time.0
1393            } else {
1394                Timestamp::now()
1395            };
1396
1397            Ok(ServiceEvent {
1398                service_id,
1399                process_id,
1400                status,
1401                restart_count,
1402                time: DateTime::from_timestamp_nanos(
1403                    time.as_nanosecond().try_into().expect("must fit"),
1404                ),
1405            })
1406        }
1407
1408        let stream = watcher(self.pod_api.clone(), self.watch_pod_params())
1409            .touched_objects()
1410            .filter_map(|object| async move {
1411                match object {
1412                    Ok(pod) => Some(into_service_event(pod)),
1413                    Err(error) => {
1414                        // We assume that errors returned by Kubernetes are usually transient, so we
1415                        // just log a warning and ignore them otherwise.
1416                        tracing::warn!("service watch error: {error}");
1417                        None
1418                    }
1419                }
1420            });
1421        Box::pin(stream)
1422    }
1423
1424    fn update_scheduling_config(&self, config: ServiceSchedulingConfig) {
1425        *self.scheduling_config.write().expect("poisoned") = config;
1426    }
1427}
1428
1429impl OrchestratorWorker {
1430    fn spawn(self, name: String) -> AbortOnDropHandle<()> {
1431        mz_ore::task::spawn(|| name, self.run()).abort_on_drop()
1432    }
1433
1434    async fn run(mut self) {
1435        {
1436            info!("initializing Kubernetes orchestrator worker");
1437            let start = Instant::now();
1438
1439            // Fetch the owner reference for our own pod (usually a
1440            // StatefulSet), so that we can propagate it to the services we
1441            // create.
1442            let hostname = env::var("HOSTNAME").unwrap_or_else(|_| panic!("HOSTNAME environment variable missing or invalid; required for Kubernetes orchestrator"));
1443            let orchestrator_pod = Retry::default()
1444                .clamp_backoff(Duration::from_secs(10))
1445                .retry_async(|_| self.pod_api.get(&hostname))
1446                .await
1447                .expect("always retries on error");
1448            self.owner_references
1449                .extend(orchestrator_pod.owner_references().into_iter().cloned());
1450
1451            if !self.collect_pod_metrics {
1452                info!(
1453                    "pod metrics collection is disabled; resource usage graphs in the console will not be available"
1454                );
1455            }
1456
1457            info!(
1458                "Kubernetes orchestrator worker initialized in {:?}",
1459                start.elapsed()
1460            );
1461        }
1462
1463        while let Some(cmd) = self.command_rx.recv().await {
1464            self.handle_command(cmd).await;
1465        }
1466    }
1467
1468    /// Handle a worker command.
1469    ///
1470    /// If handling the command fails, it is automatically retried. All command handlers return
1471    /// [`K8sError`], so we can reasonably assume that a failure is caused by issues communicating
1472    /// with the K8S server and that retrying resolves them eventually.
1473    async fn handle_command(&self, cmd: WorkerCommand) {
1474        async fn retry<F, U, R>(f: F, cmd_type: &str) -> R
1475        where
1476            F: Fn() -> U,
1477            U: Future<Output = Result<R, K8sError>>,
1478        {
1479            let start = Instant::now();
1480            Retry::default()
1481                .clamp_backoff(Duration::from_secs(10))
1482                .retry_async(|state| {
1483                    f().map_err(move |error| {
1484                        let attempt = state.i + 1;
1485                        let elapsed = start.elapsed();
1486                        if state.i < RETRY_WARN_ATTEMPTS {
1487                            tracing::warn!(
1488                                %cmd_type,
1489                                attempt,
1490                                ?elapsed,
1491                                "orchestrator call failed: {error}"
1492                            );
1493                        } else {
1494                            tracing::error!(
1495                                %cmd_type,
1496                                attempt,
1497                                ?elapsed,
1498                                "orchestrator call failed: {error}"
1499                            );
1500                        }
1501                    })
1502                })
1503                .await
1504                .expect("always retries on error")
1505        }
1506
1507        use WorkerCommand::*;
1508        match cmd {
1509            EnsureService { desc } => {
1510                retry(|| self.ensure_service(desc.clone()), "EnsureService").await
1511            }
1512            DropService { name } => retry(|| self.drop_service(&name), "DropService").await,
1513            ListServices {
1514                namespace,
1515                result_tx,
1516            } => {
1517                let result = retry(|| self.list_services(&namespace), "ListServices").await;
1518                let _ = result_tx.send(result);
1519            }
1520            Flush { result_tx } => {
1521                let _ = result_tx.send(());
1522            }
1523            FetchServiceMetrics {
1524                name,
1525                info,
1526                result_tx,
1527            } => {
1528                let result = self.fetch_service_metrics(&name, &info).await;
1529                let _ = result_tx.send(result);
1530            }
1531        }
1532    }
1533
1534    async fn fetch_service_metrics(
1535        &self,
1536        name: &str,
1537        info: &ServiceInfo,
1538    ) -> Vec<ServiceProcessMetrics> {
1539        if !self.collect_pod_metrics {
1540            return (0..info.scale.get())
1541                .map(|_| ServiceProcessMetrics::default())
1542                .collect();
1543        }
1544
1545        /// Usage metrics reported by clusterd processes.
1546        #[derive(Deserialize)]
1547        pub(crate) struct ClusterdUsage {
1548            disk_bytes: Option<u64>,
1549            memory_bytes: Option<u64>,
1550            swap_bytes: Option<u64>,
1551            heap_limit: Option<u64>,
1552        }
1553
1554        /// Get metrics for a particular service and process, converting them into a sane (i.e., numeric) format.
1555        ///
1556        /// Note that we want to keep going even if a lookup fails for whatever reason,
1557        /// so this function is infallible. If we fail to get cpu or memory for a particular pod,
1558        /// we just log a warning and install `None` in the returned struct.
1559        async fn get_metrics(
1560            self_: &OrchestratorWorker,
1561            service_name: &str,
1562            i: usize,
1563        ) -> ServiceProcessMetrics {
1564            let name = format!("{service_name}-{i}");
1565
1566            let clusterd_usage_fut = get_clusterd_usage(self_, service_name, i);
1567            let (metrics, clusterd_usage) =
1568                match futures::future::join(self_.metrics_api.get(&name), clusterd_usage_fut).await
1569                {
1570                    (Ok(metrics), Ok(clusterd_usage)) => (metrics, Some(clusterd_usage)),
1571                    (Ok(metrics), Err(e)) => {
1572                        warn!("Failed to fetch clusterd usage for {name}: {e}");
1573                        (metrics, None)
1574                    }
1575                    (Err(e), _) => {
1576                        warn!("Failed to get metrics for {name}: {e}");
1577                        return ServiceProcessMetrics::default();
1578                    }
1579                };
1580            let Some(PodMetricsContainer {
1581                usage:
1582                    PodMetricsContainerUsage {
1583                        cpu: Quantity(cpu_str),
1584                        memory: Quantity(mem_str),
1585                    },
1586                ..
1587            }) = metrics.containers.get(0)
1588            else {
1589                warn!("metrics result contained no containers for {name}");
1590                return ServiceProcessMetrics::default();
1591            };
1592
1593            let mut process_metrics = ServiceProcessMetrics::default();
1594
1595            match parse_k8s_quantity(cpu_str) {
1596                Ok(q) => match q.try_to_integer(-9, true) {
1597                    Some(nano_cores) => process_metrics.cpu_nano_cores = Some(nano_cores),
1598                    None => error!("CPU value {q:?} out of range"),
1599                },
1600                Err(e) => error!("failed to parse CPU value {cpu_str}: {e}"),
1601            }
1602            match parse_k8s_quantity(mem_str) {
1603                Ok(q) => match q.try_to_integer(0, false) {
1604                    Some(mem) => process_metrics.memory_bytes = Some(mem),
1605                    None => error!("memory value {q:?} out of range"),
1606                },
1607                Err(e) => error!("failed to parse memory value {mem_str}: {e}"),
1608            }
1609
1610            if let Some(usage) = clusterd_usage {
1611                // clusterd may report disk usage as either `disk_bytes`, or `swap_bytes`, or both.
1612                //
1613                // For now the Console expects the swap size to be reported in `disk_bytes`.
1614                // Once the Console has been ported to use `heap_bytes`/`heap_limit`, we can
1615                // simplify things by setting `process_metrics.disk_bytes = usage.disk_bytes`.
1616                process_metrics.disk_bytes = match (usage.disk_bytes, usage.swap_bytes) {
1617                    (Some(disk), Some(swap)) => Some(disk + swap),
1618                    (disk, swap) => disk.or(swap),
1619                };
1620
1621                // clusterd may report heap usage as `memory_bytes` and optionally `swap_bytes`.
1622                // If no `memory_bytes` is reported, we can't know the heap usage.
1623                process_metrics.heap_bytes = match (usage.memory_bytes, usage.swap_bytes) {
1624                    (Some(memory), Some(swap)) => Some(memory + swap),
1625                    (Some(memory), None) => Some(memory),
1626                    (None, _) => None,
1627                };
1628
1629                process_metrics.heap_limit = usage.heap_limit;
1630                process_metrics.swap_bytes = usage.swap_bytes;
1631            }
1632
1633            process_metrics
1634        }
1635
1636        /// Get the current usage metrics exposed by a clusterd process.
1637        ///
1638        /// Usage metrics are collected by connecting to a metrics endpoint exposed by the process.
1639        /// The endpoint is assumed to be reachable at the 'internal-http' under the HTTP path
1640        /// `/api/usage-metrics`.
1641        async fn get_clusterd_usage(
1642            self_: &OrchestratorWorker,
1643            service_name: &str,
1644            i: usize,
1645        ) -> anyhow::Result<ClusterdUsage> {
1646            let service = self_
1647                .service_api
1648                .get(service_name)
1649                .await
1650                .with_context(|| format!("failed to get service {service_name}"))?;
1651            let namespace = service
1652                .metadata
1653                .namespace
1654                .context("missing service namespace")?;
1655            let internal_http_port = service
1656                .spec
1657                .and_then(|spec| spec.ports)
1658                .and_then(|ports| {
1659                    ports
1660                        .into_iter()
1661                        .find(|p| p.name == Some("internal-http".into()))
1662                })
1663                .map(|p| p.port);
1664            let Some(port) = internal_http_port else {
1665                bail!("internal-http port missing in service spec");
1666            };
1667            let metrics_url = format!(
1668                "http://{service_name}-{i}.{service_name}.{namespace}.svc.cluster.local:{port}\
1669                 /api/usage-metrics"
1670            );
1671
1672            let http_client = reqwest::Client::builder()
1673                .timeout(Duration::from_secs(10))
1674                .build()
1675                .context("error building HTTP client")?;
1676            let resp = http_client.get(metrics_url).send().await?;
1677            let usage = resp.json().await?;
1678
1679            Ok(usage)
1680        }
1681
1682        let ret = futures::future::join_all(
1683            (0..info.scale.cast_into()).map(|i| get_metrics(self, name, i)),
1684        );
1685
1686        ret.await
1687    }
1688
1689    async fn ensure_service(&self, mut desc: ServiceDescription) -> Result<(), K8sError> {
1690        // We inject our own pod's owner references into the Kubernetes objects
1691        // created for the service so that if the
1692        // Deployment/StatefulSet/whatever that owns the pod running the
1693        // orchestrator gets deleted, so do all services spawned by this
1694        // orchestrator.
1695        desc.service
1696            .metadata
1697            .owner_references
1698            .get_or_insert(vec![])
1699            .extend(self.owner_references.iter().cloned());
1700        desc.stateful_set
1701            .metadata
1702            .owner_references
1703            .get_or_insert(vec![])
1704            .extend(self.owner_references.iter().cloned());
1705
1706        let ss_spec = desc.stateful_set.spec.as_ref().unwrap();
1707        let pod_metadata = ss_spec.template.metadata.as_ref().unwrap();
1708        let pod_annotations = pod_metadata.annotations.clone();
1709
1710        self.service_api
1711            .patch(
1712                &desc.name,
1713                &PatchParams::apply(FIELD_MANAGER).force(),
1714                &Patch::Apply(desc.service),
1715            )
1716            .await?;
1717        self.stateful_set_api
1718            .patch(
1719                &desc.name,
1720                &PatchParams::apply(FIELD_MANAGER).force(),
1721                &Patch::Apply(desc.stateful_set),
1722            )
1723            .await?;
1724
1725        // We manage pod recreation manually, using the OnDelete StatefulSet update strategy, for
1726        // two reasons:
1727        //  * Kubernetes doesn't always automatically replace StatefulSet pods when their specs
1728        //    change, see https://github.com/kubernetes/kubernetes#67250.
1729        //  * Kubernetes replaces StatefulSet pods when their annotations change, which is not
1730        //    something we want as it could cause unavailability.
1731        //
1732        // Our pod recreation policy is simple: If a pod's template hash changed, delete it, and
1733        // let the StatefulSet controller recreate it. Otherwise, patch the existing pod's
1734        // annotations to line up with the ones in the spec.
1735        for pod_id in 0..desc.scale.get() {
1736            let pod_name = format!("{}-{pod_id}", desc.name);
1737            let pod = match self.pod_api.get(&pod_name).await {
1738                Ok(pod) => pod,
1739                // Pod already doesn't exist.
1740                Err(kube::Error::Api(e)) if e.code == 404 => continue,
1741                Err(e) => return Err(e),
1742            };
1743
1744            let result = if pod.annotations().get(POD_TEMPLATE_HASH_ANNOTATION)
1745                != Some(&desc.pod_template_hash)
1746            {
1747                self.pod_api
1748                    .delete(&pod_name, &DeleteParams::default())
1749                    .await
1750                    .map(|_| ())
1751            } else {
1752                let metadata = ObjectMeta {
1753                    annotations: pod_annotations.clone(),
1754                    ..Default::default()
1755                }
1756                .into_request_partial::<Pod>();
1757                self.pod_api
1758                    .patch_metadata(
1759                        &pod_name,
1760                        &PatchParams::apply(FIELD_MANAGER).force(),
1761                        &Patch::Apply(&metadata),
1762                    )
1763                    .await
1764                    .map(|_| ())
1765            };
1766
1767            match result {
1768                Ok(()) => (),
1769                // Pod was deleted concurrently.
1770                Err(kube::Error::Api(e)) if e.code == 404 => continue,
1771                Err(e) => return Err(e),
1772            }
1773        }
1774
1775        Ok(())
1776    }
1777
1778    async fn drop_service(&self, name: &str) -> Result<(), K8sError> {
1779        let res = self
1780            .stateful_set_api
1781            .delete(name, &DeleteParams::default())
1782            .await;
1783        match res {
1784            Ok(_) => (),
1785            Err(K8sError::Api(e)) if e.code == 404 => (),
1786            Err(e) => return Err(e),
1787        }
1788
1789        let res = self
1790            .service_api
1791            .delete(name, &DeleteParams::default())
1792            .await;
1793        match res {
1794            Ok(_) => Ok(()),
1795            Err(K8sError::Api(e)) if e.code == 404 => Ok(()),
1796            Err(e) => Err(e),
1797        }
1798    }
1799
1800    async fn list_services(&self, namespace: &str) -> Result<Vec<String>, K8sError> {
1801        let stateful_sets = self.stateful_set_api.list(&Default::default()).await?;
1802        let name_prefix = format!("{}{namespace}-", self.name_prefix);
1803        Ok(stateful_sets
1804            .into_iter()
1805            .filter_map(|ss| {
1806                ss.metadata
1807                    .name
1808                    .unwrap()
1809                    .strip_prefix(&name_prefix)
1810                    .map(Into::into)
1811            })
1812            .collect())
1813    }
1814}
1815
1816#[derive(Debug, Clone)]
1817struct KubernetesService {
1818    hosts: Vec<String>,
1819    ports: BTreeMap<String, u16>,
1820}
1821
1822impl Service for KubernetesService {
1823    fn addresses(&self, port: &str) -> Vec<String> {
1824        let port = self.ports[port];
1825        self.hosts
1826            .iter()
1827            .map(|host| format!("{host}:{port}"))
1828            .collect()
1829    }
1830}
1831
1832/// Returns the `minDomains` value for a `TopologySpreadConstraint`.
1833///
1834/// `minDomains` must be suppressed when spread is soft (Kubernetes rejects
1835/// `minDomains` with `ScheduleAnyway`) and when `availability_zones` is set
1836/// (node affinity already constrains eligible domains; if `minDomains` exceeds
1837/// the number of pinned zones the global minimum is treated as 0, causing all
1838/// but one replica to remain pending with `maxSkew=1`).
1839fn topology_spread_min_domains(
1840    soft: bool,
1841    az_pinned: bool,
1842    min_domains: Option<i32>,
1843) -> Option<i32> {
1844    if soft || az_pinned { None } else { min_domains }
1845}
1846
1847#[cfg(test)]
1848mod tests {
1849    use super::*;
1850
1851    #[mz_ore::test]
1852    fn topology_spread_min_domains_suppression() {
1853        // min_domains is kept when neither soft nor az-pinned
1854        assert_eq!(topology_spread_min_domains(false, false, Some(3)), Some(3));
1855        // min_domains is None when not set regardless of flags
1856        assert_eq!(topology_spread_min_domains(false, false, None), None);
1857        // suppressed when soft (Kubernetes rejects minDomains with ScheduleAnyway)
1858        assert_eq!(topology_spread_min_domains(true, false, Some(3)), None);
1859        // suppressed when availability_zones pins to specific AZs
1860        assert_eq!(topology_spread_min_domains(false, true, Some(3)), None);
1861        // suppressed when both
1862        assert_eq!(topology_spread_min_domains(true, true, Some(3)), None);
1863    }
1864
1865    #[mz_ore::test]
1866    fn k8s_quantity_base10_large() {
1867        let cases = &[
1868            ("42", 42),
1869            ("42k", 42000),
1870            ("42M", 42000000),
1871            ("42G", 42000000000),
1872            ("42T", 42000000000000),
1873            ("42P", 42000000000000000),
1874        ];
1875
1876        for (input, expected) in cases {
1877            let quantity = parse_k8s_quantity(input).unwrap();
1878            let number = quantity.try_to_integer(0, true).unwrap();
1879            assert_eq!(number, *expected, "input={input}, quantity={quantity:?}");
1880        }
1881    }
1882
1883    #[mz_ore::test]
1884    fn k8s_quantity_base10_small() {
1885        let cases = &[("42n", 42), ("42u", 42000), ("42m", 42000000)];
1886
1887        for (input, expected) in cases {
1888            let quantity = parse_k8s_quantity(input).unwrap();
1889            let number = quantity.try_to_integer(-9, true).unwrap();
1890            assert_eq!(number, *expected, "input={input}, quantity={quantity:?}");
1891        }
1892    }
1893
1894    #[mz_ore::test]
1895    fn k8s_quantity_base2() {
1896        let cases = &[
1897            ("42Ki", 42 << 10),
1898            ("42Mi", 42 << 20),
1899            ("42Gi", 42 << 30),
1900            ("42Ti", 42 << 40),
1901            ("42Pi", 42 << 50),
1902        ];
1903
1904        for (input, expected) in cases {
1905            let quantity = parse_k8s_quantity(input).unwrap();
1906            let number = quantity.try_to_integer(0, false).unwrap();
1907            assert_eq!(number, *expected, "input={input}, quantity={quantity:?}");
1908        }
1909    }
1910}