1use std::collections::{BTreeMap, BTreeSet};
13use std::fmt;
14use std::num::NonZero;
15use std::str::FromStr;
16use std::sync::Arc;
17use std::sync::LazyLock;
18use std::time::Duration;
19
20use anyhow::anyhow;
21use bytesize::ByteSize;
22use chrono::{DateTime, Utc};
23use futures::stream::{BoxStream, StreamExt};
24use mz_cluster_client::client::{ClusterReplicaLocation, TimelyConfig};
25use mz_compute_client::logging::LogVariant;
26use mz_compute_types::config::{ComputeReplicaConfig, ComputeReplicaLogging};
27use mz_controller_types::dyncfgs::{
28 ARRANGEMENT_EXERT_PROPORTIONALITY, CONTROLLER_PAST_GENERATION_REPLICA_CLEANUP_RETRY_INTERVAL,
29 ENABLE_TIMELY_ZERO_COPY, ENABLE_TIMELY_ZERO_COPY_LGALLOC, ENABLE_UNIFIED_CLUSTER,
30 TIMELY_ZERO_COPY_LIMIT,
31};
32use mz_controller_types::{ClusterId, ReplicaId};
33use mz_orchestrator::NamespacedOrchestrator;
34use mz_orchestrator::{
35 CpuLimit, DiskLimit, LabelSelectionLogic, LabelSelector, MemoryLimit, Service, ServiceConfig,
36 ServiceEvent, ServicePort,
37};
38use mz_ore::cast::CastInto;
39use mz_ore::task::{self, AbortOnDropHandle};
40use mz_ore::{halt, instrument};
41use mz_repr::GlobalId;
42use mz_repr::adt::numeric::Numeric;
43use regex::Regex;
44use serde::{Deserialize, Serialize};
45use tokio::time;
46use tracing::{error, info, warn};
47
48use crate::Controller;
49
50pub struct ClusterConfig {
52 pub arranged_logs: BTreeMap<LogVariant, GlobalId>,
57 pub workload_class: Option<String>,
60}
61
62pub type ClusterStatus = mz_orchestrator::ServiceStatus;
64
65#[derive(Clone, Debug, Serialize, PartialEq)]
67pub struct ReplicaConfig {
68 pub location: ReplicaLocation,
70 pub compute: ComputeReplicaConfig,
72}
73
74#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
76pub struct ReplicaAllocation {
77 pub memory_limit: Option<MemoryLimit>,
79 pub cpu_limit: Option<CpuLimit>,
81 pub cpu_request: Option<CpuLimit>,
83 pub disk_limit: Option<DiskLimit>,
85 pub scale: NonZero<u16>,
87 pub workers: NonZero<usize>,
89 #[serde(deserialize_with = "mz_repr::adt::numeric::str_serde::deserialize")]
91 pub credits_per_hour: Numeric,
92 #[serde(default)]
94 pub cpu_exclusive: bool,
95 #[serde(default = "default_true")]
98 pub is_cc: bool,
99 #[serde(default)]
108 pub family: Option<String>,
109 #[serde(default)]
111 pub swap_enabled: bool,
112 #[serde(default)]
114 pub disabled: bool,
115 #[serde(default)]
117 pub selectors: BTreeMap<String, String>,
118}
119
120impl ReplicaAllocation {
121 pub fn family(&self) -> &str {
130 match &self.family {
131 Some(family) => family.as_str(),
132 None if self.is_cc => "cc",
133 None => "legacy",
134 }
135 }
136}
137
138fn default_true() -> bool {
139 true
140}
141
142#[mz_ore::test]
143#[cfg_attr(miri, ignore)] fn test_replica_allocation_deserialization() {
146 use bytesize::ByteSize;
147 use mz_ore::{assert_err, assert_ok};
148
149 let data = r#"
150 {
151 "cpu_limit": 1.0,
152 "memory_limit": "10GiB",
153 "disk_limit": "100MiB",
154 "scale": 16,
155 "workers": 1,
156 "credits_per_hour": "16",
157 "swap_enabled": true,
158 "selectors": {
159 "key1": "value1",
160 "key2": "value2"
161 }
162 }"#;
163
164 let replica_allocation: ReplicaAllocation = serde_json::from_str(data)
165 .expect("deserialization from JSON succeeds for ReplicaAllocation");
166
167 assert_eq!(
168 replica_allocation,
169 ReplicaAllocation {
170 credits_per_hour: 16.into(),
171 disk_limit: Some(DiskLimit(ByteSize::mib(100))),
172 disabled: false,
173 memory_limit: Some(MemoryLimit(ByteSize::gib(10))),
174 cpu_limit: Some(CpuLimit::from_millicpus(1000)),
175 cpu_request: None,
176 cpu_exclusive: false,
177 is_cc: true,
178 family: None,
179 swap_enabled: true,
180 scale: NonZero::new(16).unwrap(),
181 workers: NonZero::new(1).unwrap(),
182 selectors: BTreeMap::from([
183 ("key1".to_string(), "value1".to_string()),
184 ("key2".to_string(), "value2".to_string())
185 ]),
186 }
187 );
188
189 let data = r#"
190 {
191 "cpu_limit": 0,
192 "memory_limit": "0GiB",
193 "disk_limit": "0MiB",
194 "scale": 1,
195 "workers": 1,
196 "credits_per_hour": "0",
197 "cpu_exclusive": true,
198 "disabled": true
199 }"#;
200
201 let replica_allocation: ReplicaAllocation = serde_json::from_str(data)
202 .expect("deserialization from JSON succeeds for ReplicaAllocation");
203
204 assert_eq!(
205 replica_allocation,
206 ReplicaAllocation {
207 credits_per_hour: 0.into(),
208 disk_limit: Some(DiskLimit(ByteSize::mib(0))),
209 disabled: true,
210 memory_limit: Some(MemoryLimit(ByteSize::gib(0))),
211 cpu_limit: Some(CpuLimit::from_millicpus(0)),
212 cpu_request: None,
213 cpu_exclusive: true,
214 is_cc: true,
215 family: None,
216 swap_enabled: false,
217 scale: NonZero::new(1).unwrap(),
218 workers: NonZero::new(1).unwrap(),
219 selectors: Default::default(),
220 }
221 );
222
223 let data = r#"{"scale": 0, "workers": 1, "credits_per_hour": "0"}"#;
225 assert_err!(serde_json::from_str::<ReplicaAllocation>(data));
226 let data = r#"{"scale": 1, "workers": 0, "credits_per_hour": "0"}"#;
227 assert_err!(serde_json::from_str::<ReplicaAllocation>(data));
228 let data = r#"{"scale": 1, "workers": 1, "credits_per_hour": "0"}"#;
229 assert_ok!(serde_json::from_str::<ReplicaAllocation>(data));
230}
231
232#[mz_ore::test]
233#[cfg_attr(miri, ignore)] fn test_replica_allocation_family() {
235 let parse = |json: &str| -> ReplicaAllocation {
236 serde_json::from_str(json).expect("deserialization from JSON succeeds")
237 };
238
239 assert_eq!(
241 parse(r#"{"scale": 1, "workers": 1, "credits_per_hour": "0", "family": "D"}"#).family(),
242 "D"
243 );
244 assert_eq!(
247 parse(r#"{"scale": 1, "workers": 1, "credits_per_hour": "0"}"#).family(),
248 "cc"
249 );
250 assert_eq!(
253 parse(r#"{"scale": 1, "workers": 1, "credits_per_hour": "0", "is_cc": false}"#).family(),
254 "legacy"
255 );
256 assert_eq!(
258 parse(
259 r#"{"scale": 1, "workers": 1, "credits_per_hour": "0", "is_cc": false, "family": "legacy-special"}"#
260 )
261 .family(),
262 "legacy-special"
263 );
264}
265
266#[derive(Clone, Debug, Serialize, PartialEq)]
268pub enum ReplicaLocation {
269 Unmanaged(UnmanagedReplicaLocation),
271 Managed(ManagedReplicaLocation),
273}
274
275impl ReplicaLocation {
276 pub fn num_processes(&self) -> usize {
278 match self {
279 ReplicaLocation::Unmanaged(UnmanagedReplicaLocation {
280 computectl_addrs, ..
281 }) => computectl_addrs.len(),
282 ReplicaLocation::Managed(ManagedReplicaLocation { allocation, .. }) => {
283 allocation.scale.cast_into()
284 }
285 }
286 }
287
288 pub fn billed_as(&self) -> Option<&str> {
289 match self {
290 ReplicaLocation::Managed(ManagedReplicaLocation { billed_as, .. }) => {
291 billed_as.as_deref()
292 }
293 ReplicaLocation::Unmanaged(_) => None,
294 }
295 }
296
297 pub fn internal(&self) -> bool {
298 match self {
299 ReplicaLocation::Managed(ManagedReplicaLocation { internal, .. }) => *internal,
300 ReplicaLocation::Unmanaged(_) => false,
301 }
302 }
303
304 pub fn workers(&self) -> Option<usize> {
308 match self {
309 ReplicaLocation::Managed(ManagedReplicaLocation { allocation, .. }) => {
310 Some(allocation.workers.get() * self.num_processes())
311 }
312 ReplicaLocation::Unmanaged(_) => None,
313 }
314 }
315
316 pub fn pending(&self) -> bool {
322 match self {
323 ReplicaLocation::Managed(ManagedReplicaLocation { pending, .. }) => *pending,
324 ReplicaLocation::Unmanaged(_) => false,
325 }
326 }
327}
328
329#[derive(Debug, Clone)]
332pub enum ClusterRole {
333 SystemCritical,
336 System,
340 User,
343}
344
345#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
347pub struct UnmanagedReplicaLocation {
348 pub storagectl_addrs: Vec<String>,
351 pub computectl_addrs: Vec<String>,
354}
355
356#[derive(Clone, Debug, Serialize, PartialEq)]
358pub struct ManagedReplicaLocation {
359 pub allocation: ReplicaAllocation,
361 pub size: String,
363 pub internal: bool,
365 pub billed_as: Option<String>,
367 #[serde(skip)]
378 pub availability_zones: Vec<String>,
379 pub pending: bool,
381}
382
383impl ManagedReplicaLocation {
384 pub fn size_for_billing(&self) -> &str {
386 self.billed_as.as_deref().unwrap_or(&self.size)
387 }
388}
389
390pub type ReplicaLogging = ComputeReplicaLogging;
392
393pub type ProcessId = u64;
395
396#[derive(Debug, Clone, Serialize)]
398pub struct ClusterEvent {
399 pub cluster_id: ClusterId,
400 pub replica_id: ReplicaId,
401 pub process_id: ProcessId,
402 pub status: ClusterStatus,
403 pub restart_count: u64,
406 pub time: DateTime<Utc>,
407}
408
409impl Controller {
410 pub fn create_cluster(
416 &mut self,
417 id: ClusterId,
418 config: ClusterConfig,
419 ) -> Result<(), anyhow::Error> {
420 self.storage
421 .create_instance(id, config.workload_class.clone());
422 self.compute
423 .create_instance(id, config.arranged_logs, config.workload_class)?;
424 Ok(())
425 }
426
427 pub fn update_cluster_workload_class(&mut self, id: ClusterId, workload_class: Option<String>) {
433 self.storage
434 .update_instance_workload_class(id, workload_class.clone());
435 self.compute
436 .update_instance_workload_class(id, workload_class)
437 .expect("instance exists");
438 }
439
440 pub fn drop_cluster(&mut self, id: ClusterId) {
446 self.storage.drop_instance(id);
447 self.compute.drop_instance(id);
448 }
449
450 pub fn create_replica(
453 &mut self,
454 cluster_id: ClusterId,
455 replica_id: ReplicaId,
456 cluster_name: String,
457 replica_name: String,
458 role: ClusterRole,
459 config: ReplicaConfig,
460 enable_worker_core_affinity: bool,
461 ) -> Result<(), anyhow::Error> {
462 let storage_location: ClusterReplicaLocation;
463 let compute_location: ClusterReplicaLocation;
464 let metrics_task: Option<AbortOnDropHandle<()>>;
465
466 match config.location {
467 ReplicaLocation::Unmanaged(UnmanagedReplicaLocation {
468 storagectl_addrs,
469 computectl_addrs,
470 }) => {
471 compute_location = ClusterReplicaLocation {
472 ctl_addrs: computectl_addrs,
473 };
474 storage_location = ClusterReplicaLocation {
475 ctl_addrs: storagectl_addrs,
476 };
477 metrics_task = None;
478 }
479 ReplicaLocation::Managed(m) => {
480 let (service, metrics_task_join_handle) = self.provision_replica(
481 cluster_id,
482 replica_id,
483 cluster_name,
484 replica_name,
485 role,
486 m,
487 enable_worker_core_affinity,
488 )?;
489 storage_location = ClusterReplicaLocation {
490 ctl_addrs: service.addresses("storagectl"),
491 };
492 compute_location = ClusterReplicaLocation {
493 ctl_addrs: service.addresses("computectl"),
494 };
495 metrics_task = Some(metrics_task_join_handle);
496
497 let http_addresses = service.addresses("internal-http");
499 self.replica_http_locator
500 .register_replica(cluster_id, replica_id, http_addresses);
501 }
502 }
503
504 self.storage
505 .connect_replica(cluster_id, replica_id, storage_location);
506 self.compute.add_replica_to_instance(
507 cluster_id,
508 replica_id,
509 compute_location,
510 config.compute,
511 )?;
512
513 if let Some(task) = metrics_task {
514 self.metrics_tasks.insert(replica_id, task);
515 }
516
517 Ok(())
518 }
519
520 pub fn drop_replica(
522 &mut self,
523 cluster_id: ClusterId,
524 replica_id: ReplicaId,
525 ) -> Result<(), anyhow::Error> {
526 self.deprovision_replica(cluster_id, replica_id, self.deploy_generation)?;
531 self.metrics_tasks.remove(&replica_id);
532
533 self.replica_http_locator
535 .remove_replica(cluster_id, replica_id);
536
537 self.replica_dyncfg_overrides.remove(&replica_id);
541
542 self.compute.drop_replica(cluster_id, replica_id)?;
543 self.storage.drop_replica(cluster_id, replica_id);
544 Ok(())
545 }
546
547 pub(crate) fn remove_past_generation_replicas_in_background(&self) {
549 let deploy_generation = self.deploy_generation;
550 let dyncfg = Arc::clone(self.compute.dyncfg());
551 let orchestrator = Arc::clone(&self.orchestrator);
552 task::spawn(
553 || "controller_remove_past_generation_replicas",
554 async move {
555 info!("attempting to remove past generation replicas");
556 loop {
557 match try_remove_past_generation_replicas(&*orchestrator, deploy_generation)
558 .await
559 {
560 Ok(()) => {
561 info!("successfully removed past generation replicas");
562 return;
563 }
564 Err(e) => {
565 let interval =
566 CONTROLLER_PAST_GENERATION_REPLICA_CLEANUP_RETRY_INTERVAL
567 .get(&dyncfg);
568 warn!(%e, "failed to remove past generation replicas; will retry in {interval:?}");
569 time::sleep(interval).await;
570 }
571 }
572 }
573 },
574 );
575 }
576
577 #[instrument]
579 pub async fn remove_orphaned_replicas(
580 &mut self,
581 next_user_replica_id: u64,
582 next_system_replica_id: u64,
583 ) -> Result<(), anyhow::Error> {
584 let desired: BTreeSet<_> = self.metrics_tasks.keys().copied().collect();
585
586 let actual: BTreeSet<_> = self
587 .orchestrator
588 .list_services()
589 .await?
590 .iter()
591 .map(|s| ReplicaServiceName::from_str(s))
592 .collect::<Result<_, _>>()?;
593
594 for ReplicaServiceName {
595 cluster_id,
596 replica_id,
597 generation,
598 } in actual
599 {
600 if generation != self.deploy_generation {
604 continue;
605 }
606
607 let smaller_next = match replica_id {
608 ReplicaId::User(id) if id >= next_user_replica_id => {
609 Some(ReplicaId::User(next_user_replica_id))
610 }
611 ReplicaId::System(id) if id >= next_system_replica_id => {
612 Some(ReplicaId::System(next_system_replica_id))
613 }
614 _ => None,
615 };
616 if let Some(next) = smaller_next {
617 halt!("found replica ID ({replica_id}) in orchestrator >= next ID ({next})");
622 }
623 if !desired.contains(&replica_id) {
624 self.deprovision_replica(cluster_id, replica_id, generation)?;
625 }
626 }
627
628 self.orchestrator.flush().await?;
629 Ok(())
630 }
631
632 pub fn events_stream(&self) -> BoxStream<'static, ClusterEvent> {
633 let deploy_generation = self.deploy_generation;
634
635 fn translate_event(event: ServiceEvent) -> Result<(ClusterEvent, u64), anyhow::Error> {
636 let ReplicaServiceName {
637 cluster_id,
638 replica_id,
639 generation: replica_generation,
640 ..
641 } = event.service_id.parse()?;
642
643 let event = ClusterEvent {
644 cluster_id,
645 replica_id,
646 process_id: event.process_id,
647 status: event.status,
648 restart_count: event.restart_count,
649 time: event.time,
650 };
651
652 Ok((event, replica_generation))
653 }
654
655 let stream = self
656 .orchestrator
657 .watch_services()
658 .map(|event| event.and_then(translate_event))
659 .filter_map(move |event| async move {
660 match event {
661 Ok((event, replica_generation)) => {
662 if replica_generation == deploy_generation {
663 Some(event)
664 } else {
665 None
666 }
667 }
668 Err(error) => {
669 error!("service watch error: {error}");
670 None
671 }
672 }
673 });
674
675 Box::pin(stream)
676 }
677
678 fn provision_replica(
680 &self,
681 cluster_id: ClusterId,
682 replica_id: ReplicaId,
683 cluster_name: String,
684 replica_name: String,
685 role: ClusterRole,
686 location: ManagedReplicaLocation,
687 enable_worker_core_affinity: bool,
688 ) -> Result<(Box<dyn Service>, AbortOnDropHandle<()>), anyhow::Error> {
689 let service_name = ReplicaServiceName {
690 cluster_id,
691 replica_id,
692 generation: self.deploy_generation,
693 }
694 .to_string();
695 let role_label = match role {
696 ClusterRole::SystemCritical => "system-critical",
697 ClusterRole::System => "system",
698 ClusterRole::User => "user",
699 };
700 let environment_id = self.connection_context().environment_id.clone();
701 let aws_external_id_prefix = self.connection_context().aws_external_id_prefix.clone();
702 let aws_connection_role_arn = self.connection_context().aws_connection_role_arn.clone();
703 let persist_pubsub_url = self.persist_pubsub_url.clone();
704 let secrets_args = self.secrets_args.to_flags();
705
706 let overrides = self.replica_dyncfg_overrides.get(&replica_id);
713 let arrangement_exert_proportionality =
716 ARRANGEMENT_EXERT_PROPORTIONALITY.get_with_overrides(&self.dyncfg, overrides);
717 let storage_proto_timely_config = TimelyConfig {
718 arrangement_exert_proportionality,
719 ..Default::default()
720 };
721 let compute_proto_timely_config = TimelyConfig {
722 arrangement_exert_proportionality,
723 enable_zero_copy: ENABLE_TIMELY_ZERO_COPY.get_with_overrides(&self.dyncfg, overrides),
724 enable_zero_copy_lgalloc: ENABLE_TIMELY_ZERO_COPY_LGALLOC
725 .get_with_overrides(&self.dyncfg, overrides),
726 zero_copy_limit: TIMELY_ZERO_COPY_LIMIT.get_with_overrides(&self.dyncfg, overrides),
727 ..Default::default()
728 };
729 let unified_cluster = ENABLE_UNIFIED_CLUSTER.get_with_overrides(&self.dyncfg, overrides);
730
731 let mut disk_limit = location.allocation.disk_limit;
732 let memory_limit = location.allocation.memory_limit;
733 let mut memory_request = None;
734
735 if location.allocation.swap_enabled {
736 disk_limit = Some(DiskLimit::ZERO);
740
741 memory_request = memory_limit.map(|MemoryLimit(limit)| {
745 let request = ByteSize::b(limit.as_u64() - 1);
746 MemoryLimit(request)
747 });
748 }
749
750 let service = self.orchestrator.ensure_service(
751 &service_name,
752 ServiceConfig {
753 app_name: "clusterd".into(),
754 image: self.clusterd_image.clone(),
755 init_container_image: self.init_container_image.clone(),
756 args: Box::new(move |assigned| {
757 let storage_timely_config = TimelyConfig {
758 workers: location.allocation.workers.get(),
759 addresses: assigned.peer_addresses("storage"),
760 ..storage_proto_timely_config
761 };
762 let compute_timely_config = TimelyConfig {
763 workers: location.allocation.workers.get(),
764 addresses: assigned.peer_addresses("compute"),
765 ..compute_proto_timely_config
766 };
767
768 let mut args = vec![
769 format!(
770 "--storage-controller-listen-addr={}",
771 assigned.listen_addrs["storagectl"]
772 ),
773 format!(
774 "--compute-controller-listen-addr={}",
775 assigned.listen_addrs["computectl"]
776 ),
777 format!(
778 "--internal-http-listen-addr={}",
779 assigned.listen_addrs["internal-http"]
780 ),
781 format!("--opentelemetry-resource=cluster_id={}", cluster_id),
782 format!("--opentelemetry-resource=replica_id={}", replica_id),
783 format!("--persist-pubsub-url={}", persist_pubsub_url),
784 format!("--environment-id={}", environment_id),
785 format!(
786 "--storage-timely-config={}",
787 storage_timely_config.to_string(),
788 ),
789 format!(
790 "--compute-timely-config={}",
791 compute_timely_config.to_string(),
792 ),
793 ];
794 if let Some(aws_external_id_prefix) = &aws_external_id_prefix {
795 args.push(format!(
796 "--aws-external-id-prefix={}",
797 aws_external_id_prefix
798 ));
799 }
800 if let Some(aws_connection_role_arn) = &aws_connection_role_arn {
801 args.push(format!(
802 "--aws-connection-role-arn={}",
803 aws_connection_role_arn
804 ));
805 }
806 if let Some(memory_limit) = location.allocation.memory_limit {
807 args.push(format!(
808 "--announce-memory-limit={}",
809 memory_limit.0.as_u64()
810 ));
811 }
812 if location.allocation.cpu_exclusive && enable_worker_core_affinity {
813 args.push("--worker-core-affinity".into());
814 }
815 if unified_cluster {
816 args.push("--unified-cluster".into());
817 }
818 if location.allocation.is_cc {
819 args.push("--is-cc".into());
820 }
821
822 if location.allocation.swap_enabled
825 && let Some(memory_limit) = location.allocation.memory_limit
826 && let Some(disk_limit) = location.allocation.disk_limit
827 && disk_limit != DiskLimit::ZERO
831 {
832 let heap_limit = memory_limit.0 + disk_limit.0;
833 args.push(format!("--heap-limit={}", heap_limit.as_u64()));
834 }
835
836 args.extend(secrets_args.clone());
837 args
838 }),
839 ports: vec![
840 ServicePort {
841 name: "storagectl".into(),
842 port_hint: 2100,
843 },
844 ServicePort {
848 name: "storage".into(),
849 port_hint: 2103,
850 },
851 ServicePort {
852 name: "computectl".into(),
853 port_hint: 2101,
854 },
855 ServicePort {
856 name: "compute".into(),
857 port_hint: 2102,
858 },
859 ServicePort {
860 name: "internal-http".into(),
861 port_hint: 6878,
862 },
863 ],
864 cpu_limit: location.allocation.cpu_limit,
865 cpu_request: location.allocation.cpu_request,
866 memory_limit,
867 memory_request,
868 scale: location.allocation.scale,
869 labels: BTreeMap::from([
870 ("replica-id".into(), replica_id.to_string()),
871 ("cluster-id".into(), cluster_id.to_string()),
872 ("generation".into(), self.deploy_generation.to_string()),
873 ("type".into(), "cluster".into()),
874 ("replica-role".into(), role_label.into()),
875 ("workers".into(), location.allocation.workers.to_string()),
876 (
877 "size".into(),
878 location
879 .size
880 .to_string()
881 .replace("=", "-")
882 .replace(",", "_"),
883 ),
884 ]),
885 annotations: BTreeMap::from([
886 (
887 "replica-name".into(),
888 format!("{cluster_name}.{replica_name}"),
889 ),
890 ("cluster-name".into(), cluster_name),
891 ]),
892 availability_zones: Some(location.availability_zones).filter(|azs| !azs.is_empty()),
895 other_replicas_selector: vec![
908 LabelSelector {
909 label_name: "cluster-id".to_string(),
910 logic: LabelSelectionLogic::Eq {
911 value: cluster_id.to_string(),
912 },
913 },
914 LabelSelector {
916 label_name: "replica-id".into(),
917 logic: LabelSelectionLogic::NotEq {
918 value: replica_id.to_string(),
919 },
920 },
921 LabelSelector {
922 label_name: "generation".into(),
923 logic: LabelSelectionLogic::Eq {
924 value: self.deploy_generation.to_string(),
925 },
926 },
927 ],
928 replicas_selector: vec![
929 LabelSelector {
930 label_name: "cluster-id".to_string(),
931 logic: LabelSelectionLogic::Eq {
933 value: cluster_id.to_string(),
934 },
935 },
936 LabelSelector {
937 label_name: "generation".into(),
938 logic: LabelSelectionLogic::Eq {
939 value: self.deploy_generation.to_string(),
940 },
941 },
942 ],
943 disk_limit,
944 node_selector: location.allocation.selectors,
945 },
946 )?;
947
948 let metrics_task = mz_ore::task::spawn(|| format!("replica-metrics-{replica_id}"), {
949 let tx = self.metrics_tx.clone();
950 let orchestrator = Arc::clone(&self.orchestrator);
951 let service_name = service_name.clone();
952 async move {
953 const METRICS_INTERVAL: Duration = Duration::from_secs(60);
954
955 let mut interval = tokio::time::interval(METRICS_INTERVAL);
963 loop {
964 interval.tick().await;
965 match orchestrator.fetch_service_metrics(&service_name).await {
966 Ok(metrics) => {
967 let _ = tx.send((replica_id, metrics));
968 }
969 Err(e) => {
970 warn!("failed to get metrics for replica {replica_id}: {e}");
971 }
972 }
973 }
974 }
975 });
976
977 Ok((service, metrics_task.abort_on_drop()))
978 }
979
980 fn deprovision_replica(
982 &self,
983 cluster_id: ClusterId,
984 replica_id: ReplicaId,
985 generation: u64,
986 ) -> Result<(), anyhow::Error> {
987 let service_name = ReplicaServiceName {
988 cluster_id,
989 replica_id,
990 generation,
991 }
992 .to_string();
993 self.orchestrator.drop_service(&service_name)
994 }
995}
996
997async fn try_remove_past_generation_replicas(
999 orchestrator: &dyn NamespacedOrchestrator,
1000 deploy_generation: u64,
1001) -> Result<(), anyhow::Error> {
1002 let services: BTreeSet<_> = orchestrator.list_services().await?.into_iter().collect();
1003
1004 for service in services {
1005 let name: ReplicaServiceName = service.parse()?;
1006 if name.generation < deploy_generation {
1007 info!(
1008 cluster_id = %name.cluster_id,
1009 replica_id = %name.replica_id,
1010 "removing past generation replica",
1011 );
1012 orchestrator.drop_service(&service)?;
1013 }
1014 }
1015
1016 Ok(())
1017}
1018
1019#[derive(PartialEq, Eq, PartialOrd, Ord)]
1021pub struct ReplicaServiceName {
1022 pub cluster_id: ClusterId,
1023 pub replica_id: ReplicaId,
1024 pub generation: u64,
1025}
1026
1027impl fmt::Display for ReplicaServiceName {
1028 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
1029 let ReplicaServiceName {
1030 cluster_id,
1031 replica_id,
1032 generation,
1033 } = self;
1034 write!(f, "{cluster_id}-replica-{replica_id}-gen-{generation}")
1035 }
1036}
1037
1038impl FromStr for ReplicaServiceName {
1039 type Err = anyhow::Error;
1040
1041 fn from_str(s: &str) -> Result<Self, Self::Err> {
1042 static SERVICE_NAME_RE: LazyLock<Regex> = LazyLock::new(|| {
1043 Regex::new(r"(?-u)^([us]\d+)-replica-([us]\d+)(?:-gen-(\d+))?$").unwrap()
1044 });
1045
1046 let caps = SERVICE_NAME_RE
1047 .captures(s)
1048 .ok_or_else(|| anyhow!("invalid service name: {s}"))?;
1049
1050 Ok(ReplicaServiceName {
1051 cluster_id: caps.get(1).unwrap().as_str().parse().unwrap(),
1052 replica_id: caps.get(2).unwrap().as_str().parse().unwrap(),
1053 generation: caps.get(3).map_or("0", |m| m.as_str()).parse().unwrap(),
1057 })
1058 }
1059}