Skip to main content

mz_orchestrator/
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::fmt;
12use std::num::NonZero;
13use std::str::FromStr;
14use std::sync::Arc;
15
16use async_trait::async_trait;
17use bytesize::ByteSize;
18use chrono::{DateTime, Utc};
19use derivative::Derivative;
20use futures_core::stream::BoxStream;
21use mz_ore::cast::CastFrom;
22use serde::de::Unexpected;
23use serde::{Deserialize, Deserializer, Serialize};
24
25/// An orchestrator manages services.
26///
27/// A service is a set of one or more processes running the same image. See
28/// [`ServiceConfig`] for details.
29///
30/// All services live within a namespace. A namespace allows multiple users to
31/// share an orchestrator without conflicting: each user can only create,
32/// delete, and list the services within their namespace. Namespaces are not
33/// isolated at the network level, however: services in one namespace can
34/// communicate with services in another namespace with no restrictions.
35///
36/// Services **must** be tolerant of running as part of a distributed system. In
37/// particular, services **must** be prepared for the possibility that there are
38/// two live processes with the same identity. This can happen, for example,
39/// when the machine hosting a process *appears* to fail, from the perspective
40/// of the orchestrator, and so the orchestrator restarts the process on another
41/// machine, but in fact the original machine is still alive, just on the
42/// opposite side of a network partition. Be sure to design any communication
43/// with other services (e.g., an external database) to correctly handle
44/// competing communication from another incarnation of the service.
45///
46/// The intent is that you can implement `Orchestrator` with pods in Kubernetes,
47/// containers in Docker, or processes on your local machine.
48pub trait Orchestrator: fmt::Debug + Send + Sync {
49    /// Enter a namespace in the orchestrator.
50    fn namespace(&self, namespace: &str) -> Arc<dyn NamespacedOrchestrator>;
51}
52
53/// An orchestrator restricted to a single namespace.
54#[async_trait]
55pub trait NamespacedOrchestrator: fmt::Debug + Send + Sync {
56    /// Ensures that a service with the given configuration is running.
57    ///
58    /// If a service with the same ID already exists, its configuration is
59    /// updated to match `config`. This may or may not involve restarting the
60    /// service, depending on whether the existing service matches `config`.
61    fn ensure_service(
62        &self,
63        id: &str,
64        config: ServiceConfig,
65    ) -> Result<Box<dyn Service>, anyhow::Error>;
66
67    /// Drops the identified service, if it exists.
68    fn drop_service(&self, id: &str) -> Result<(), anyhow::Error>;
69
70    /// Lists the identifiers of all known services.
71    async fn list_services(&self) -> Result<Vec<String>, anyhow::Error>;
72
73    /// Waits for previously queued requests to be processed, not for affected
74    /// processes to start or terminate.
75    async fn flush(&self) -> Result<(), anyhow::Error>;
76
77    /// Watch for status changes of all known services.
78    fn watch_services(&self) -> BoxStream<'static, Result<ServiceEvent, anyhow::Error>>;
79
80    /// Gets resource usage metrics for all processes associated with a service.
81    ///
82    /// Returns `Err` if the entire process failed. Returns `Ok(v)` otherwise,
83    /// with one element in `v` for each process of the service,
84    /// even in not all metrics could be collected for all processes.
85    /// In such a case, the corresponding fields of `ServiceProcessMetrics` will be `None`.
86    async fn fetch_service_metrics(
87        &self,
88        id: &str,
89    ) -> Result<Vec<ServiceProcessMetrics>, anyhow::Error>;
90
91    fn update_scheduling_config(&self, config: scheduling_config::ServiceSchedulingConfig);
92}
93
94/// An event describing a status change of an orchestrated service.
95#[derive(Debug, Clone, Serialize)]
96pub struct ServiceEvent {
97    pub service_id: String,
98    pub process_id: u64,
99    pub status: ServiceStatus,
100    /// Cumulative number of times the underlying process has restarted, as
101    /// reported by the orchestrator. Monotonic for the lifetime of a process,
102    /// but can reset (e.g. when a pod is recreated). Orchestrators that don't
103    /// track restarts report 0.
104    pub restart_count: u64,
105    pub time: DateTime<Utc>,
106}
107
108/// Why the service is not ready, if known
109#[derive(Debug, Clone, Copy, Serialize, Eq, PartialEq)]
110pub enum OfflineReason {
111    OomKilled,
112    Initializing,
113}
114
115impl fmt::Display for OfflineReason {
116    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
117        match self {
118            OfflineReason::OomKilled => f.write_str("oom-killed"),
119            OfflineReason::Initializing => f.write_str("initializing"),
120        }
121    }
122}
123
124/// Describes the status of an orchestrated service.
125#[derive(Debug, Clone, Copy, Serialize, Eq, PartialEq)]
126pub enum ServiceStatus {
127    /// Service is ready to accept requests.
128    Online,
129    /// Service is not ready to accept requests.
130    /// The inner element is `None` if the reason
131    /// is unknown
132    Offline(Option<OfflineReason>),
133}
134
135impl ServiceStatus {
136    /// Returns the service status as a kebab-case string.
137    pub fn as_kebab_case_str(&self) -> &'static str {
138        match self {
139            ServiceStatus::Online => "online",
140            ServiceStatus::Offline(_) => "offline",
141        }
142    }
143}
144
145/// Describes a running service managed by an `Orchestrator`.
146pub trait Service: fmt::Debug + Send + Sync {
147    /// Given the name of a port, returns the addresses for each of the
148    /// service's processes, in order.
149    ///
150    /// Panics if `port` does not name a valid port.
151    fn addresses(&self, port: &str) -> Vec<String>;
152}
153
154#[derive(Copy, Clone, Debug, Default, Serialize, Deserialize, Eq, PartialEq)]
155pub struct ServiceProcessMetrics {
156    pub cpu_nano_cores: Option<u64>,
157    pub memory_bytes: Option<u64>,
158    pub disk_bytes: Option<u64>,
159    pub heap_bytes: Option<u64>,
160    pub heap_limit: Option<u64>,
161    pub swap_bytes: Option<u64>,
162}
163
164/// A simple language for describing assertions about a label's existence and value.
165///
166/// Used by [`LabelSelector`].
167#[derive(Clone, Debug)]
168pub enum LabelSelectionLogic {
169    /// The label exists and its value equals the given value.
170    /// Equivalent to `InSet { values: vec![value] }`
171    Eq { value: String },
172    /// Either the label does not exist, or it exists
173    /// but its value does not equal the given value.
174    /// Equivalent to `NotInSet { values: vec![value] }`
175    NotEq { value: String },
176    /// The label exists.
177    Exists,
178    /// The label does not exist.
179    NotExists,
180    /// The label exists and its value is one of the given values.
181    InSet { values: Vec<String> },
182    /// Either the label does not exist, or it exists
183    /// but its value is not one of the given values.
184    NotInSet { values: Vec<String> },
185}
186
187/// A simple language for describing whether a label
188/// exists and whether the value corresponding to it is in some set.
189/// Intended to correspond to the capabilities offered by Kubernetes label selectors,
190/// but without directly exposing Kubernetes API code to consumers of this module.
191#[derive(Clone, Debug)]
192pub struct LabelSelector {
193    /// The name of the label
194    pub label_name: String,
195    /// An assertion about the existence and value of a label
196    /// named `label_name`
197    pub logic: LabelSelectionLogic,
198}
199
200/// Describes the desired state of a service.
201#[derive(Derivative)]
202#[derivative(Debug)]
203pub struct ServiceConfig {
204    /// Static application name (usually present in labels)
205    pub app_name: String,
206    /// An opaque identifier for the executable or container image to run.
207    ///
208    /// Often names a container on Docker Hub or a path on the local machine.
209    pub image: String,
210    /// For the Kubernetes orchestrator, this is an init container to
211    /// configure for the pod running the service.
212    pub init_container_image: Option<String>,
213    /// A function that generates the arguments for each process of the service
214    /// given the assigned listen addresses for each named port.
215    #[derivative(Debug = "ignore")]
216    pub args: Box<dyn Fn(ServiceAssignments) -> Vec<String> + Send + Sync>,
217    /// Ports to expose.
218    pub ports: Vec<ServicePort>,
219    /// An optional limit on the memory that the service can use.
220    pub memory_limit: Option<MemoryLimit>,
221    /// An optional request on the memory that the service can use. If unspecified,
222    /// use the same value as `memory_limit`.
223    pub memory_request: Option<MemoryLimit>,
224    /// An optional limit on the CPU that the service can use.
225    pub cpu_limit: Option<CpuLimit>,
226    /// An optional request on the CPU that the service can use.
227    pub cpu_request: Option<CpuLimit>,
228    /// The number of copies of this service to run.
229    pub scale: NonZero<u16>,
230    /// Arbitrary key–value pairs to attach to the service in the orchestrator
231    /// backend.
232    ///
233    /// The orchestrator backend may apply a prefix to the key if appropriate.
234    pub labels: BTreeMap<String, String>,
235    /// Arbitrary key–value pairs to attach to the service as annotations in the
236    /// orchestrator backend.
237    ///
238    /// The orchestrator backend may apply a prefix to the key if appropriate.
239    pub annotations: BTreeMap<String, String>,
240    /// The availability zones the service can be run in. If no availability
241    /// zones are specified, the orchestrator is free to choose one.
242    pub availability_zones: Option<Vec<String>>,
243    /// A set of label selectors selecting all _other_ services that are replicas of this one.
244    ///
245    /// This may be used to implement anti-affinity. If _all_ such selectors
246    /// match for a given service, this service should not be co-scheduled on
247    /// a machine with that service.
248    ///
249    /// The orchestrator backend may or may not actually implement anti-affinity functionality.
250    pub other_replicas_selector: Vec<LabelSelector>,
251    /// A set of label selectors selecting all services that are replicas of this one,
252    /// including itself.
253    ///
254    /// This may be used to implement placement spread.
255    ///
256    /// The orchestrator backend may or may not actually implement placement spread functionality.
257    pub replicas_selector: Vec<LabelSelector>,
258
259    /// The maximum amount of scratch disk space that the service is allowed to consume.
260    pub disk_limit: Option<DiskLimit>,
261    /// Node selector for this service.
262    pub node_selector: BTreeMap<String, String>,
263}
264
265/// Get the recommended Kubernetes labels (app.kubernetes.io/*)
266/// WARNING: this is duplicated in src/orchestratord/src/k8s.rs and src/cloud-resources/src/crd.rs
267pub fn recommended_k8s_labels(app_name: String) -> BTreeMap<String, String> {
268    BTreeMap::from_iter([
269        (
270            "app.kubernetes.io/managed-by".to_owned(),
271            "materialize-operator".to_owned(),
272        ),
273        (
274            "app.kubernetes.io/part-of".to_owned(),
275            "materialize".to_owned(),
276        ),
277        ("app.kubernetes.io/name".to_owned(), app_name.to_owned()),
278        // legacy label
279        ("app".to_owned(), app_name.to_owned()),
280    ])
281}
282
283/// A named port associated with a service.
284#[derive(Debug, Clone, PartialEq, Eq)]
285pub struct ServicePort {
286    /// A descriptive name for the port.
287    ///
288    /// Note that not all orchestrator backends make use of port names.
289    pub name: String,
290    /// The desired port number.
291    ///
292    /// Not all orchestrator backends will make use of the hint.
293    pub port_hint: u16,
294}
295
296/// Assignments that the orchestrator has made for a process in a service.
297#[derive(Clone, Debug)]
298pub struct ServiceAssignments<'a> {
299    /// For each specified [`ServicePort`] name, a listen address.
300    pub listen_addrs: &'a BTreeMap<String, String>,
301    /// The listen addresses of each peer in the service.
302    ///
303    /// The order of peers is significant. Each peer is uniquely identified by its position in the
304    /// list.
305    pub peer_addrs: &'a [BTreeMap<String, String>],
306}
307
308impl ServiceAssignments<'_> {
309    /// Return the peer addresses for the specified [`ServicePort`] name.
310    pub fn peer_addresses(&self, name: &str) -> Vec<String> {
311        self.peer_addrs.iter().map(|a| a[name].clone()).collect()
312    }
313}
314
315/// Describes a limit on memory.
316#[derive(Copy, Clone, Debug, PartialOrd, Eq, Ord, PartialEq)]
317pub struct MemoryLimit(pub ByteSize);
318
319impl MemoryLimit {
320    pub const MAX: Self = Self(ByteSize(u64::MAX));
321}
322
323impl<'de> Deserialize<'de> for MemoryLimit {
324    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
325    where
326        D: Deserializer<'de>,
327    {
328        <String as Deserialize>::deserialize(deserializer)
329            .and_then(|s| {
330                ByteSize::from_str(&s).map_err(|_e| {
331                    use serde::de::Error;
332                    D::Error::invalid_value(serde::de::Unexpected::Str(&s), &"valid size in bytes")
333                })
334            })
335            .map(MemoryLimit)
336    }
337}
338
339impl Serialize for MemoryLimit {
340    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
341    where
342        S: serde::Serializer,
343    {
344        <String as Serialize>::serialize(&self.0.to_string(), serializer)
345    }
346}
347
348/// Describes a limit on CPU resources.
349#[derive(Debug, Copy, Clone, Eq, Ord, PartialEq, PartialOrd)]
350pub struct CpuLimit {
351    millicpus: usize,
352}
353
354impl CpuLimit {
355    pub const MAX: Self = Self::from_millicpus(usize::MAX / 1_000_000);
356
357    /// Constructs a new CPU limit from a number of millicpus.
358    pub const fn from_millicpus(millicpus: usize) -> CpuLimit {
359        CpuLimit { millicpus }
360    }
361
362    /// Returns the CPU limit in millicpus.
363    pub fn as_millicpus(&self) -> usize {
364        self.millicpus
365    }
366
367    /// Returns the CPU limit in nanocpus.
368    pub fn as_nanocpus(&self) -> u64 {
369        // The largest possible value of a u64 is
370        // 18_446_744_073_709_551_615,
371        // so we won't overflow this
372        // unless we have an instance with
373        // ~18.45 billion cores.
374        //
375        // Such an instance seems unrealistic,
376        // at least until we raise another few rounds
377        // of funding ...
378
379        u64::cast_from(self.millicpus)
380            .checked_mul(1_000_000)
381            .expect("Nano-CPUs must be representable")
382    }
383}
384
385impl<'de> Deserialize<'de> for CpuLimit {
386    // TODO(benesch): remove this once this function no longer makes use of
387    // potentially dangerous `as` conversions.
388    #[allow(clippy::as_conversions)]
389    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
390    where
391        D: serde::Deserializer<'de>,
392    {
393        // Note -- we just round off any precision beyond 0.001 here.
394        let float = f64::deserialize(deserializer)?;
395        let millicpus = (float * 1000.).round();
396        if millicpus < 0. || millicpus > (usize::MAX as f64) {
397            use serde::de::Error;
398            Err(D::Error::invalid_value(
399                Unexpected::Float(float),
400                &"a float representing a plausible number of CPUs",
401            ))
402        } else {
403            Ok(Self {
404                millicpus: millicpus as usize,
405            })
406        }
407    }
408}
409
410impl Serialize for CpuLimit {
411    // TODO(benesch): remove this once this function no longer makes use of
412    // potentially dangerous `as` conversions.
413    #[allow(clippy::as_conversions)]
414    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
415    where
416        S: serde::Serializer,
417    {
418        <f64 as Serialize>::serialize(&(self.millicpus as f64 / 1000.0), serializer)
419    }
420}
421
422/// Describes a limit on disk usage.
423#[derive(Copy, Clone, Debug, PartialOrd, Eq, Ord, PartialEq)]
424pub struct DiskLimit(pub ByteSize);
425
426impl DiskLimit {
427    pub const ZERO: Self = Self(ByteSize(0));
428    pub const MAX: Self = Self(ByteSize(u64::MAX));
429    pub const ARBITRARY: Self = Self(ByteSize::gib(1));
430}
431
432impl<'de> Deserialize<'de> for DiskLimit {
433    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
434    where
435        D: Deserializer<'de>,
436    {
437        <String as Deserialize>::deserialize(deserializer)
438            .and_then(|s| {
439                ByteSize::from_str(&s).map_err(|_e| {
440                    use serde::de::Error;
441                    D::Error::invalid_value(serde::de::Unexpected::Str(&s), &"valid size in bytes")
442                })
443            })
444            .map(DiskLimit)
445    }
446}
447
448impl Serialize for DiskLimit {
449    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
450    where
451        S: serde::Serializer,
452    {
453        <String as Serialize>::serialize(&self.0.to_string(), serializer)
454    }
455}
456
457/// Configuration for how services are scheduled. These may be ignored by orchestrator
458/// implementations.
459pub mod scheduling_config {
460    #[derive(Debug, Clone)]
461    pub struct ServiceTopologySpreadConfig {
462        /// If `true`, enable spread for replicated services.
463        ///
464        /// Defaults to `true`.
465        pub enabled: bool,
466        /// If `true`, ignore services with `scale` > 1 when expressing
467        /// spread constraints.
468        ///
469        /// Default to `true`.
470        pub ignore_non_singular_scale: bool,
471        /// The `maxSkew` for spread constraints.
472        /// See
473        /// <https://kubernetes.io/docs/concepts/scheduling-eviction/topology-spread-constraints/>
474        /// for more details.
475        ///
476        /// Defaults to `1`.
477        pub max_skew: i32,
478        /// The `minDomains` for spread constraints.
479        /// See
480        /// <https://kubernetes.io/docs/concepts/scheduling-eviction/topology-spread-constraints/>
481        /// for more details.
482        ///
483        /// Defaults to None.
484        pub min_domains: Option<i32>,
485        /// If `true`, make the spread constraints into a preference.
486        ///
487        /// Defaults to `false`.
488        pub soft: bool,
489    }
490
491    #[derive(Debug, Clone)]
492    pub struct ServiceSchedulingConfig {
493        /// If `Some`, add a affinity preference with the given
494        /// weight for services that horizontally scale.
495        ///
496        /// Defaults to `Some(100)`.
497        pub multi_pod_az_affinity_weight: Option<i32>,
498        /// If `true`, make the node-scope anti-affinity between
499        /// replicated services a preference over a constraint.
500        ///
501        /// Defaults to `false`.
502        pub soften_replication_anti_affinity: bool,
503        /// The weight for `soften_replication_anti_affinity.
504        ///
505        /// Defaults to `100`.
506        pub soften_replication_anti_affinity_weight: i32,
507        /// Configuration for `TopologySpreadConstraint`'s
508        pub topology_spread: ServiceTopologySpreadConfig,
509        /// If `true`, make the az-scope node affinity soft.
510        ///
511        /// Defaults to `false`.
512        pub soften_az_affinity: bool,
513        /// The weight for `soften_replication_anti_affinity.
514        ///
515        /// Defaults to `100`.
516        pub soften_az_affinity_weight: i32,
517        // Whether to enable security context for the service.
518        pub security_context_enabled: bool,
519    }
520
521    pub const DEFAULT_POD_AZ_AFFINITY_WEIGHT: Option<i32> = Some(100);
522    pub const DEFAULT_SOFTEN_REPLICATION_ANTI_AFFINITY: bool = false;
523    pub const DEFAULT_SOFTEN_REPLICATION_ANTI_AFFINITY_WEIGHT: i32 = 100;
524
525    pub const DEFAULT_TOPOLOGY_SPREAD_ENABLED: bool = true;
526    pub const DEFAULT_TOPOLOGY_SPREAD_IGNORE_NON_SINGULAR_SCALE: bool = true;
527    pub const DEFAULT_TOPOLOGY_SPREAD_MAX_SKEW: i32 = 1;
528    pub const DEFAULT_TOPOLOGY_SPREAD_MIN_DOMAIN: Option<i32> = None;
529    pub const DEFAULT_TOPOLOGY_SPREAD_SOFT: bool = false;
530
531    pub const DEFAULT_SOFTEN_AZ_AFFINITY: bool = false;
532    pub const DEFAULT_SOFTEN_AZ_AFFINITY_WEIGHT: i32 = 100;
533    pub const DEFAULT_SECURITY_CONTEXT_ENABLED: bool = true;
534
535    impl Default for ServiceSchedulingConfig {
536        fn default() -> Self {
537            ServiceSchedulingConfig {
538                multi_pod_az_affinity_weight: DEFAULT_POD_AZ_AFFINITY_WEIGHT,
539                soften_replication_anti_affinity: DEFAULT_SOFTEN_REPLICATION_ANTI_AFFINITY,
540                soften_replication_anti_affinity_weight:
541                    DEFAULT_SOFTEN_REPLICATION_ANTI_AFFINITY_WEIGHT,
542                topology_spread: ServiceTopologySpreadConfig {
543                    enabled: DEFAULT_TOPOLOGY_SPREAD_ENABLED,
544                    ignore_non_singular_scale: DEFAULT_TOPOLOGY_SPREAD_IGNORE_NON_SINGULAR_SCALE,
545                    max_skew: DEFAULT_TOPOLOGY_SPREAD_MAX_SKEW,
546                    min_domains: DEFAULT_TOPOLOGY_SPREAD_MIN_DOMAIN,
547                    soft: DEFAULT_TOPOLOGY_SPREAD_SOFT,
548                },
549                soften_az_affinity: DEFAULT_SOFTEN_AZ_AFFINITY,
550                soften_az_affinity_weight: DEFAULT_SOFTEN_AZ_AFFINITY_WEIGHT,
551                security_context_enabled: DEFAULT_SECURITY_CONTEXT_ENABLED,
552            }
553        }
554    }
555}