1use 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;
68const RETRY_WARN_ATTEMPTS: usize = 5;
73
74const POD_TEMPLATE_HASH_ANNOTATION: &str = "environmentd.materialize.cloud/pod-template-hash";
75
76#[derive(Debug, Clone)]
78pub struct KubernetesOrchestratorConfig {
79 pub context: String,
82 pub scheduler_name: Option<String>,
84 pub priority_class_name: Option<String>,
86 pub service_annotations: BTreeMap<String, String>,
88 pub service_labels: BTreeMap<String, String>,
90 pub service_node_selector: BTreeMap<String, String>,
92 pub service_affinity: Option<String>,
94 pub service_tolerations: Option<String>,
96 pub service_account: Option<String>,
98 pub image_pull_policy: KubernetesImagePullPolicy,
100 pub aws_external_id_prefix: Option<AwsExternalIdPrefix>,
103 pub coverage: bool,
105 pub ephemeral_volume_storage_class: Option<String>,
111 pub service_fs_group: Option<i64>,
113 pub name_prefix: Option<String>,
115 pub collect_pod_metrics: bool,
117 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#[derive(ValueEnum, Debug, Clone, Copy)]
129pub enum KubernetesImagePullPolicy {
130 Always,
132 IfNotPresent,
134 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
158pub 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 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 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
253enum 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#[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
289struct 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#[derive(Deserialize, Clone, Debug)]
362pub struct MetricIdentifier {
363 #[serde(rename = "metricName")]
364 pub name: String,
365 }
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 }
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 fn watch_pod_params(&self) -> watcher::Config {
391 let ns_selector = format!(
392 "environmentd.materialize.cloud/namespace={}",
393 self.namespace
394 );
395 watcher::Config::default().timeout(59).labels(&ns_selector)
397 }
398
399 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
474fn 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), ("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 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 let scheduling_config: ServiceSchedulingConfig =
590 self.scheduling_config.read().expect("poisoned").clone();
591
592 let disk = disk_limit != Some(DiskLimit::ZERO);
594
595 let name = self.service_name(id);
596 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 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 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 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 match_label_keys: None,
843 ..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 "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 "karpenter.sh/do-not-disrupt".to_owned() => "true".to_string(),
866 };
867 for (key, value) in annotations_in {
868 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 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 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 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 ..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 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 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 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 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 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 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 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 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 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 #[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 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 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 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 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 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 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 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 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
1832fn 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 assert_eq!(topology_spread_min_domains(false, false, Some(3)), Some(3));
1855 assert_eq!(topology_spread_min_domains(false, false, None), None);
1857 assert_eq!(topology_spread_min_domains(true, false, Some(3)), None);
1859 assert_eq!(topology_spread_min_domains(false, true, Some(3)), None);
1861 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}