Skip to main content

mz_adapter/coord/
cluster_controller.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! Driver and glue for the [`mz_cluster_controller`] reconciler.
11//!
12//! The controller crate is pure: it knows nothing about the Coordinator. This
13//! module is the half of the [`ClusterControllerCtx`] boundary that does: it runs
14//! the controller as a **separate task** and implements the ctx by marshaling
15//! each pull/apply to the Coordinator over the internal command channel, because
16//! the catalog and the live compute/storage signals are reachable only from the
17//! coordinator loop. Whole-tick reads are batched. Refresh-window catalog inputs
18//! are pulled one cluster at a time and completed with one shared oracle read.
19//! The remaining per-cluster live signals are pulled on demand, so steady
20//! clusters do not pay for signals they do not use.
21//!
22//! The controller owns the replica set of every managed cluster, user and
23//! system alike. A builtin cluster's config-implied replicas are additionally
24//! materialized by `reconcile_builtin_cluster_replicas` at catalog open, which
25//! derives the same target from the same config, so the two converge rather
26//! than compete.
27
28use std::collections::{BTreeMap, BTreeSet};
29use std::sync::Arc;
30use std::time::Duration;
31
32use mz_adapter_types::dyncfgs::{
33    CLUSTER_CONTROLLER_TICK_INTERVAL, CLUSTER_RECONFIGURATION_ALLOWED_LAG,
34    ENABLE_CLUSTER_RECONFIGURATION_LAG_GATE,
35};
36use mz_catalog::memory::objects::{ClusterConfig, ClusterVariant};
37use mz_cluster_controller::ClusterController;
38use mz_cluster_controller::ctx::{
39    ApplyOutcome, AvailabilityZones, ClusterControllerCtx, ClusterState, CreateReason, Decision,
40    ExpectedClusterState, ObservedReplica, OnTimeout, ReconfigurationRecord, ReconfigurationStatus,
41    ReconfigurationTarget, RefreshMvInfo, RefreshWindowClusterInputs, RefreshWindowInputsBatch,
42    ReplicaShape, StateWrite,
43};
44use mz_cluster_controller::strategy::duration_to_ts;
45use mz_compute_types::config::ComputeReplicaConfig;
46use mz_controller::clusters::ClusterStatus;
47use mz_controller_types::{ClusterId, ReplicaId};
48use mz_ore::task::spawn;
49use mz_repr::Timestamp;
50use tokio::sync::{mpsc, oneshot};
51use tracing::{debug, warn};
52
53use crate::catalog::{DropObjectInfo, Op, ReplicaCreateDropReason};
54use crate::coord::{ClusterReplicaStatuses, Coordinator, Message};
55use crate::error::AdapterError;
56
57/// A request the controller task marshals to the Coordinator to satisfy one
58/// [`ClusterControllerCtx`] call. Each variant carries a oneshot for the reply.
59///
60/// `ManagedClusterIds` and `ClusterStates` are the per-tick batched reads. The
61/// `ClusterStates` reply also carries `now`. Refresh-window catalog inputs are
62/// pulled one cluster at a time, followed by one shared oracle read.
63/// `HydratedReplicas` and `ReadyReplicas` are per-cluster live signals a strategy
64/// pulls on demand.
65#[derive(Debug)]
66pub enum ClusterControllerRequest {
67    /// The ids of all managed clusters the controller owns this tick.
68    ManagedClusterIds { tx: oneshot::Sender<Vec<ClusterId>> },
69    /// A consistent durable view of the given clusters and their replicas, plus
70    /// the current time.
71    ClusterStates {
72        clusters: Vec<ClusterId>,
73        tx: oneshot::Sender<(Vec<ClusterState>, Timestamp)>,
74    },
75    /// Of `replicas` on `cluster`, which are online and have all current
76    /// collections hydrated.
77    HydratedReplicas {
78        cluster_id: ClusterId,
79        replicas: Vec<ReplicaId>,
80        tx: oneshot::Sender<BTreeSet<ReplicaId>>,
81    },
82    /// Of `replicas` on `cluster`, which are online, have all current
83    /// collections hydrated, *and* are within the configured cut-over lag of
84    /// the `reference` replicas, the ones the cut-over would drop.
85    ///
86    /// The stronger form of [`Self::HydratedReplicas`], for callers deciding
87    /// whether to cut over to a replica rather than merely observing that its
88    /// dataflows have started producing output.
89    ReadyReplicas {
90        cluster_id: ClusterId,
91        replicas: Vec<ReplicaId>,
92        reference: BTreeSet<ReplicaId>,
93        tx: oneshot::Sender<BTreeSet<ReplicaId>>,
94    },
95    /// Whether the cluster has any hydratable (dataflow-backed) objects bound to
96    /// it.
97    HasHydratableObjects {
98        cluster_id: ClusterId,
99        tx: oneshot::Sender<bool>,
100    },
101    /// The catalog and storage refresh-window inputs for one scheduled cluster.
102    /// `None` if the cluster no longer qualifies at pull time.
103    RefreshWindowClusterInputs {
104        cluster_id: ClusterId,
105        tx: oneshot::Sender<Option<RefreshWindowClusterInputs>>,
106    },
107    /// One timestamp-oracle read for a completed refresh-window input batch.
108    RefreshWindowReadTs { tx: oneshot::Sender<Timestamp> },
109    /// Apply a tick's batch of decisions under their compare-and-append guards.
110    Apply {
111        decisions: Vec<Decision>,
112        tx: oneshot::Sender<ApplyOutcome>,
113    },
114    /// The current configured reconcile cadence. Read once per tick so a runtime
115    /// change to `cluster_controller_tick_interval` takes effect without a
116    /// restart.
117    TickInterval { tx: oneshot::Sender<Duration> },
118}
119
120struct ReplicaReadinessCheck {
121    replica_id: ReplicaId,
122    compute_ready: oneshot::Receiver<bool>,
123}
124
125/// The controller-task side of the boundary: a [`ClusterControllerCtx`] that
126/// marshals every call to the Coordinator over `internal_cmd_tx`.
127struct CoordCtx {
128    internal_cmd_tx: mpsc::UnboundedSender<Message>,
129    /// Latched `now` from the most recent batched read, returned by
130    /// [`ClusterControllerCtx::now`] so a strategy and the kernel see a single
131    /// consistent time per phase.
132    now: Timestamp,
133}
134
135impl CoordCtx {
136    /// Send a request and await its reply. Returns `None` if the Coordinator has
137    /// gone away (shutdown), which the caller treats as "nothing to do".
138    async fn request<T>(
139        &self,
140        make: impl FnOnce(oneshot::Sender<T>) -> ClusterControllerRequest,
141    ) -> Option<T> {
142        let (tx, rx) = oneshot::channel();
143        if self
144            .internal_cmd_tx
145            .send(Message::ClusterControllerRequest(make(tx)))
146            .is_err()
147        {
148            return None;
149        }
150        rx.await.ok()
151    }
152}
153
154#[async_trait::async_trait]
155impl ClusterControllerCtx for CoordCtx {
156    fn now(&self) -> Timestamp {
157        self.now
158    }
159
160    async fn managed_cluster_ids(&mut self) -> Vec<ClusterId> {
161        self.request(|tx| ClusterControllerRequest::ManagedClusterIds { tx })
162            .await
163            .unwrap_or_default()
164    }
165
166    async fn cluster_states(&mut self, clusters: &[ClusterId]) -> Vec<ClusterState> {
167        let clusters = clusters.to_vec();
168        match self
169            .request(|tx| ClusterControllerRequest::ClusterStates { clusters, tx })
170            .await
171        {
172            Some((states, now)) => {
173                self.now = now;
174                states
175            }
176            None => Vec::new(),
177        }
178    }
179
180    async fn hydrated_replicas(
181        &mut self,
182        cluster_id: ClusterId,
183        replicas: &[ReplicaId],
184    ) -> BTreeSet<ReplicaId> {
185        let replicas = replicas.to_vec();
186        self.request(|tx| ClusterControllerRequest::HydratedReplicas {
187            cluster_id,
188            replicas,
189            tx,
190        })
191        .await
192        .unwrap_or_default()
193    }
194
195    async fn ready_replicas(
196        &mut self,
197        cluster_id: ClusterId,
198        replicas: &[ReplicaId],
199        reference: &BTreeSet<ReplicaId>,
200    ) -> BTreeSet<ReplicaId> {
201        let replicas = replicas.to_vec();
202        let reference = reference.clone();
203        self.request(|tx| ClusterControllerRequest::ReadyReplicas {
204            cluster_id,
205            replicas,
206            reference,
207            tx,
208        })
209        .await
210        .unwrap_or_default()
211    }
212
213    async fn has_hydratable_objects(&mut self, cluster_id: ClusterId) -> bool {
214        self.request(|tx| ClusterControllerRequest::HasHydratableObjects { cluster_id, tx })
215            .await
216            // A lost reply means shutdown; "no objects" arms nothing, which is
217            // the safe answer.
218            .unwrap_or(false)
219    }
220
221    async fn refresh_window_inputs(
222        &mut self,
223        cluster_ids: &[ClusterId],
224    ) -> Option<RefreshWindowInputsBatch> {
225        let mut cluster_inputs = BTreeMap::new();
226        for (index, &cluster_id) in cluster_ids.iter().enumerate() {
227            if index > 0 {
228                // The coordinator prioritizes its internal command channel. Give
229                // it a chance to service already-queued user commands instead of
230                // keeping that channel continuously ready for the whole batch.
231                tokio::task::yield_now().await;
232            }
233            let inputs = self
234                .request(|tx| ClusterControllerRequest::RefreshWindowClusterInputs {
235                    cluster_id,
236                    tx,
237                })
238                .await
239                .flatten();
240            if let Some(inputs) = inputs {
241                cluster_inputs.insert(cluster_id, inputs);
242            }
243        }
244        if cluster_inputs.is_empty() {
245            return None;
246        }
247
248        let read_ts = self
249            .request(|tx| ClusterControllerRequest::RefreshWindowReadTs { tx })
250            .await?;
251        Some(RefreshWindowInputsBatch {
252            read_ts,
253            cluster_inputs,
254        })
255    }
256
257    async fn apply(&mut self, decisions: Vec<Decision>) -> ApplyOutcome {
258        self.request(|tx| ClusterControllerRequest::Apply { decisions, tx })
259            .await
260            // A lost reply means shutdown; treat as rejected so we make no
261            // further claims about the catalog state.
262            .unwrap_or(ApplyOutcome::Rejected)
263    }
264}
265
266impl Coordinator {
267    /// Spawn the cluster controller task.
268    ///
269    /// The task ticks at [`CLUSTER_CONTROLLER_TICK_INTERVAL`], re-read each tick
270    /// via a [`ClusterControllerRequest::TickInterval`] round-trip so a runtime
271    /// change takes effect without a restart. It owns the controller and a
272    /// [`CoordCtx`] that marshals back to this Coordinator.
273    ///
274    /// The interval is the fallback cadence: `reconcile_now` cuts the
275    /// sleep short after a catalog transaction changes durable cluster state.
276    /// The notification only wakes the task. The tick still pulls fresh state
277    /// through the coordinator loop, and the controller's own applies wake it
278    /// again at the cost of one no-op tick.
279    pub(crate) fn spawn_cluster_controller_task(&self) {
280        let internal_cmd_tx = self.internal_cmd_tx.clone();
281        let reconcile_now = Arc::clone(&self.reconcile_now);
282        // A shared handle: dyncfg updates land in the same underlying values, so
283        // the controller task sees flag flips without any push.
284        let dyncfgs = self.catalog().system_config().dyncfgs().clone();
285
286        spawn(|| "cluster_controller", async move {
287            let controller = ClusterController::new(dyncfgs);
288            let mut ctx = CoordCtx {
289                internal_cmd_tx,
290                now: Timestamp::MIN,
291            };
292
293            loop {
294                // Re-read the cadence each tick so a runtime change takes effect.
295                // A lost reply means the Coordinator is gone; stop ticking.
296                let Some(interval) = ctx
297                    .request(|tx| ClusterControllerRequest::TickInterval { tx })
298                    .await
299                else {
300                    break;
301                };
302                tokio::select! {
303                    _ = tokio::time::sleep(interval.max(Duration::from_millis(1))) => {}
304                    _ = reconcile_now.notified() => {}
305                }
306
307                if ctx.internal_cmd_tx.is_closed() {
308                    // Coordinator gone; stop ticking.
309                    break;
310                }
311                controller.reconcile(&mut ctx).await;
312            }
313        });
314    }
315
316    /// Handle one [`ClusterControllerRequest`] on the coordinator loop.
317    ///
318    /// The controller is inactive while the deployment is in read-only mode (a
319    /// 0dt upgrade, where it must not write the catalog). When inactive, reads
320    /// report no managed clusters (so the controller finds nothing to
321    /// reconcile) and applies are rejected: the task still wakes each tick and
322    /// sends one `ManagedClusterIds` request, but that request early-returns
323    /// here and no catalog state is read or written. The task keeps ticking, so
324    /// the controller reactivates on its own once the deployment promotes out of
325    /// read-only mode.
326    #[mz_ore::instrument(level = "debug")]
327    pub(crate) async fn handle_cluster_controller_request(
328        &mut self,
329        request: ClusterControllerRequest,
330    ) {
331        let active = !self.controller.read_only();
332
333        match request {
334            ClusterControllerRequest::ManagedClusterIds { tx } => {
335                let ids = if active {
336                    self.catalog()
337                        .clusters()
338                        .filter(|c| c.is_managed())
339                        .map(|c| c.id)
340                        .collect()
341                } else {
342                    Vec::new()
343                };
344                let _ = tx.send(ids);
345            }
346            ClusterControllerRequest::ClusterStates { clusters, tx } => {
347                let now = Timestamp::from(self.now());
348                // Only ever asked about clusters the controller is reconciling
349                // this tick, which the inactive `ManagedClusterIds` gate above
350                // makes empty, so no guard is needed here.
351                let states: Vec<_> = clusters
352                    .into_iter()
353                    .filter_map(|id| self.observe_cluster_state(id))
354                    .collect();
355                let _ = tx.send((states, now));
356            }
357            ClusterControllerRequest::HydratedReplicas {
358                cluster_id,
359                replicas,
360                tx,
361            } => {
362                // Hydration only: this signal drives the burst strategy's
363                // linger, whose durable `steady_hydrated_at` stamp means exactly
364                // "hydration was observed". The cut-over gate is `ReadyReplicas`.
365                let checks =
366                    self.start_readiness_checks(cluster_id, replicas, None, &BTreeSet::new());
367                Self::finish_readiness_checks(checks, tx);
368            }
369            ClusterControllerRequest::ReadyReplicas {
370                cluster_id,
371                replicas,
372                reference,
373                tx,
374            } => {
375                // Tests hold cut-over while arranging a hydrated but lagging
376                // replacement, then release this to exercise the real probe.
377                fail::fail_point!("cluster_controller_hold_readiness", |_| ());
378                let checks = self.start_readiness_checks(
379                    cluster_id,
380                    replicas,
381                    self.reconfiguration_allowed_lag(),
382                    &reference,
383                );
384                Self::finish_readiness_checks(checks, tx);
385            }
386            ClusterControllerRequest::HasHydratableObjects { cluster_id, tx } => {
387                let _ = tx.send(self.cluster_has_hydratable_objects(cluster_id));
388            }
389            ClusterControllerRequest::RefreshWindowClusterInputs { cluster_id, tx } => {
390                let _ = tx.send(self.refresh_window_catalog_inputs(cluster_id));
391            }
392            ClusterControllerRequest::RefreshWindowReadTs { tx } => {
393                // The oracle read is a network round trip to the Postgres or
394                // CRDB-backed timestamp oracle. It must not run on the serial
395                // coordinator loop.
396                let oracle = self.get_local_timestamp_oracle();
397                spawn(|| "cluster_controller_refresh_window_read_ts", async move {
398                    let read_ts = oracle.read_ts().await;
399                    let _ = tx.send(read_ts);
400                });
401            }
402            ClusterControllerRequest::Apply { decisions, tx } => {
403                let outcome = if active {
404                    self.apply_cluster_decisions(decisions).await
405                } else {
406                    ApplyOutcome::Rejected
407                };
408                let _ = tx.send(outcome);
409            }
410            ClusterControllerRequest::TickInterval { tx } => {
411                let interval =
412                    CLUSTER_CONTROLLER_TICK_INTERVAL.get(self.catalog().system_config().dyncfgs());
413                let _ = tx.send(interval);
414            }
415        }
416    }
417
418    /// Build the controller's view of one managed cluster from the catalog.
419    /// Returns `None` for a missing or unmanaged cluster.
420    fn observe_cluster_state(&self, cluster_id: ClusterId) -> Option<ClusterState> {
421        let cluster = self.catalog().try_get_cluster(cluster_id)?;
422        let ClusterVariant::Managed(managed) = &cluster.config.variant else {
423            return None;
424        };
425        // The witness fields come from the same projection the compare-and-append
426        // check uses, so the state a decision is derived from and the state the
427        // apply path checks against cannot drift.
428        let expected = crate::catalog::cluster_state::project_expected(managed);
429
430        // All replicas, with the raw traits the controller's ownership test
431        // (`ObservedReplica::owned_shape`) classifies on.
432        let replicas = cluster
433            .replicas()
434            .map(|replica| ObservedReplica {
435                replica_id: replica.replica_id,
436                name: replica.name.clone(),
437                shape: replica_shape(&replica.config),
438                internal: replica.config.location.internal(),
439                billed_as: replica.config.location.billed_as().is_some(),
440                pending: replica.config.location.pending(),
441            })
442            .collect();
443
444        Some(ClusterState {
445            cluster_id,
446            size: expected.size,
447            replication_factor: expected.replication_factor,
448            availability_zones: expected.availability_zones.0,
449            logging: expected.logging,
450            arrangement_compression: expected.arrangement_compression,
451            schedule: expected.schedule,
452            auto_scaling_policy: expected.auto_scaling_policy,
453            reconfiguration: expected.reconfiguration,
454            burst: expected.burst,
455            replicas,
456        })
457    }
458
459    /// Whether the cluster has any hydratable objects bound to it, backing the
460    /// controller's [`ClusterControllerCtx::has_hydratable_objects`] pull (see
461    /// the trait method for the approximation contract and why mismatches with
462    /// the hydration check are self-healing).
463    fn cluster_has_hydratable_objects(&self, cluster_id: ClusterId) -> bool {
464        let Some(cluster) = self.catalog().try_get_cluster(cluster_id) else {
465            return false;
466        };
467        cluster
468            .bound_objects
469            .iter()
470            .any(|id| self.catalog().get_entry(id).item().is_hydratable())
471    }
472
473    /// The cut-over lag allowance, or `None` when the lag gate is disabled.
474    ///
475    /// Read per probe rather than latched per tick, so a runtime change takes
476    /// effect without a restart, matching `cluster_controller_tick_interval`.
477    pub(crate) fn reconfiguration_allowed_lag(&self) -> Option<Timestamp> {
478        let dyncfgs = self.catalog().system_config().dyncfgs();
479        ENABLE_CLUSTER_RECONFIGURATION_LAG_GATE
480            .get(dyncfgs)
481            // A lag beyond the timestamp domain means "any lag is fine", which
482            // is what `duration_to_ts` saturating at `MAX` expresses.
483            .then(|| duration_to_ts(CLUSTER_RECONFIGURATION_ALLOWED_LAG.get(dyncfgs)))
484    }
485
486    /// Await the started checks off the coordinator loop and reply with the
487    /// replicas that passed. The compute check can wait on the compute instance
488    /// task, so it must not block the loop.
489    fn finish_readiness_checks(
490        checks: Vec<ReplicaReadinessCheck>,
491        tx: oneshot::Sender<BTreeSet<ReplicaId>>,
492    ) {
493        spawn(|| "cluster_controller_readiness_probe", async move {
494            let mut ready = BTreeSet::new();
495            for check in checks {
496                if check.compute_ready.await.unwrap_or(false) {
497                    ready.insert(check.replica_id);
498                }
499            }
500            let _ = tx.send(ready);
501        });
502    }
503
504    /// Starts per-replica readiness checks for `cluster_id`: hydration, plus
505    /// the lag gate against `reference` when `allowed_lag` is `Some`.
506    ///
507    /// Returns only checks for replicas whose processes are all online, that
508    /// are already storage-hydrated, and that are known to the compute
509    /// controller. The compute receiver completes off the coordinator loop.
510    ///
511    /// The storage-side check is hydration only. Storage hydration has its own
512    /// definition (see `StorageController::collections_hydrated_on_replicas`),
513    /// and no lag term is applied to it here; a pending replica hosting an
514    /// ingestion can pass this gate with the source's snapshot complete but its
515    /// replay still in progress.
516    fn start_readiness_checks(
517        &self,
518        cluster_id: ClusterId,
519        replicas: Vec<ReplicaId>,
520        allowed_lag: Option<Timestamp>,
521        reference: &BTreeSet<ReplicaId>,
522    ) -> Vec<ReplicaReadinessCheck> {
523        use mz_catalog::memory::objects::CatalogItem;
524
525        // Materialized views pinned to a replica (via `IN CLUSTER ... REPLICA`)
526        // are only ever installed on that replica, so any other replica can
527        // never report them hydrated. Collect the cluster's pinned MVs once,
528        // then exclude the ones pinned elsewhere from each replica's hydration
529        // check. Otherwise a graceful reconfiguration's cut-over to a fresh
530        // replica set would wait forever for the new replicas to hydrate an MV
531        // bound to a replica being replaced, then roll back at the deadline
532        // (leaving the old replica, and the targeted MV, in place). Indexes
533        // cannot be replica-pinned, so MVs are the only case.
534        let pinned_mvs: Vec<(ReplicaId, mz_repr::GlobalId)> = self
535            .catalog()
536            .try_get_cluster(cluster_id)
537            .into_iter()
538            .flat_map(|cluster| cluster.bound_objects.iter())
539            .filter_map(|id| match self.catalog().get_entry(id).item() {
540                CatalogItem::MaterializedView(mv) => mv
541                    .target_replica
542                    .map(|target| (target, mv.global_id_writes())),
543                _ => None,
544            })
545            .collect();
546
547        let mut checks = Vec::new();
548        for replica_id in replicas {
549            // Skip replicas that are not online. We wait for a replica to be
550            // online even when it has no objects that need hydration on it
551            // (e.g. a single-replica source).
552            let status = self
553                .cluster_replica_statuses
554                .try_get_cluster_replica_statuses(cluster_id, replica_id)
555                .map(ClusterReplicaStatuses::cluster_replica_status);
556            if !matches!(status, Some(ClusterStatus::Online)) {
557                continue;
558            }
559            let exclude: BTreeSet<mz_repr::GlobalId> = pinned_mvs
560                .iter()
561                .filter(|(target, _)| *target != replica_id)
562                .map(|(_, id)| *id)
563                .collect();
564            let compute_fut = match self.controller.compute.collections_ready_for_replicas(
565                cluster_id,
566                vec![replica_id],
567                exclude.clone(),
568                allowed_lag,
569                reference.clone(),
570            ) {
571                Ok(fut) => fut,
572                // The replica is not known to the compute controller. Treat it
573                // as not ready.
574                Err(_) => continue,
575            };
576            let storage_hydrated = match self.controller.storage.collections_hydrated_on_replicas(
577                Some(vec![replica_id]),
578                &cluster_id,
579                &exclude,
580            ) {
581                Ok(hydrated) => hydrated,
582                Err(_) => continue,
583            };
584            if storage_hydrated {
585                checks.push(ReplicaReadinessCheck {
586                    replica_id,
587                    compute_ready: compute_fut,
588                });
589            }
590        }
591        checks
592    }
593
594    /// The catalog- and storage-derived refresh-window signals for one scheduled
595    /// cluster (the system compaction estimate and each bound REFRESH
596    /// materialized view's storage write frontier and refresh schedule), or
597    /// `None` if the cluster is missing, unmanaged, or not scheduled `ON
598    /// REFRESH`.
599    ///
600    /// The oracle read timestamp completing [`RefreshWindowInputsBatch`] is
601    /// deliberately not fetched here: this runs on the coordinator loop, and
602    /// the oracle read is a network round-trip the request handler performs on
603    /// a spawned task instead.
604    ///
605    /// The MV write frontier is carried through with full fidelity as the
606    /// `Antichain` the storage controller reports. The on-refresh strategy
607    /// compares against it directly.
608    fn refresh_window_catalog_inputs(
609        &self,
610        cluster_id: ClusterId,
611    ) -> Option<RefreshWindowClusterInputs> {
612        use mz_catalog::memory::objects::CatalogItem;
613
614        let cluster = self.catalog().try_get_cluster(cluster_id)?;
615        let ClusterVariant::Managed(managed) = &cluster.config.variant else {
616            return None;
617        };
618        if !matches!(
619            managed.schedule,
620            mz_sql::plan::ClusterSchedule::Refresh { .. }
621        ) {
622            return None;
623        }
624
625        let refresh_mvs = cluster
626            .bound_objects
627            .iter()
628            .filter_map(|id| {
629                let CatalogItem::MaterializedView(mv) = self.catalog().get_entry(id).item() else {
630                    return None;
631                };
632                let refresh_schedule = mv.refresh_schedule.clone()?;
633                // The storage controller knows about every MV in the catalog. The
634                // write frontier is passed through with full fidelity as the
635                // `Antichain` reported here.
636                let (_since, write_frontier) = self
637                    .controller
638                    .storage
639                    .collection_frontiers(mv.global_id_writes())
640                    .expect("storage controller knows about catalog MVs");
641                Some(RefreshMvInfo {
642                    id: mv.global_id_writes(),
643                    write_frontier,
644                    refresh_schedule,
645                })
646            })
647            .collect();
648
649        let compaction_estimate = self
650            .catalog()
651            .system_config()
652            .cluster_refresh_mv_compaction_estimate();
653
654        Some(RefreshWindowClusterInputs {
655            compaction_estimate,
656            refresh_mvs,
657        })
658    }
659
660    /// Apply one batch of decisions under their compare-and-append guards.
661    ///
662    /// The kernel calls this once per tick phase: a phase-1 batch is all
663    /// `UpdateClusterState`, a phase-2 batch is all create/drop. Either batch may
664    /// in principle be mixed; this handles both. The work is staged across four
665    /// steps: collect the per-cluster guards, pre-allocate the ids the creates
666    /// need, build the mutation ops, then commit ops and guards in one
667    /// transaction (see [`Self::commit_with_checks`] for why the guard holds).
668    /// Any step that finds the batch incoherent rejects it, and the controller
669    /// recomputes next tick.
670    async fn apply_cluster_decisions(&mut self, decisions: Vec<Decision>) -> ApplyOutcome {
671        let checks = Self::partition_checks(&decisions);
672
673        // Pre-allocate replica ids before the apply transaction (each allocation
674        // is its own durable commit, so it cannot happen inside the transaction).
675        let Some(replica_ids) = self.allocate_replica_ids_for_creates(&decisions).await else {
676            return ApplyOutcome::Rejected;
677        };
678
679        let Some(mutations) = self.build_mutation_ops(decisions, replica_ids) else {
680            return ApplyOutcome::Rejected;
681        };
682        if mutations.is_empty() {
683            // Nothing to apply, so the checks guard nothing. Skip the transaction
684            // rather than commit a check-only batch, which would still cost a
685            // durable round-trip.
686            return ApplyOutcome::Applied;
687        }
688
689        self.commit_with_checks(checks, mutations).await
690    }
691
692    /// The compare-and-append guards for a decision batch: one
693    /// `(cluster_id, expected)` per distinct cluster, in first-seen order. All
694    /// of a cluster's decisions in a tick come from one snapshot, so they share
695    /// one `expected` witness and one guard covers them all.
696    fn partition_checks(decisions: &[Decision]) -> Vec<(ClusterId, ExpectedClusterState)> {
697        let mut checks: Vec<(ClusterId, ExpectedClusterState)> = Vec::new();
698        let mut seen_clusters = BTreeSet::new();
699        for decision in decisions {
700            let (cluster_id, expected) = match decision {
701                Decision::CreateReplica {
702                    cluster_id,
703                    expected,
704                    ..
705                }
706                | Decision::DropReplica {
707                    cluster_id,
708                    expected,
709                    ..
710                }
711                | Decision::UpdateClusterState {
712                    cluster_id,
713                    expected,
714                    ..
715                } => (*cluster_id, expected),
716            };
717            if seen_clusters.insert(cluster_id) {
718                checks.push((cluster_id, expected.clone()));
719            } else {
720                debug_assert!(
721                    checks
722                        .iter()
723                        .any(|(c, e)| *c == cluster_id && e == expected),
724                    "decisions for a cluster in one tick must share one expected witness",
725                );
726            }
727        }
728        checks
729    }
730
731    /// Pre-allocate one replica id per `CreateReplica` decision, in the order the
732    /// creates appear (which is the order [`Self::build_mutation_ops`] consumes
733    /// them). Returns `None` if any allocation fails, which rejects the batch.
734    ///
735    /// `Op::CreateClusterReplica` carries a pre-allocated id, so we allocate
736    /// out-of-band here, before the apply transaction. Each allocation commits
737    /// durably, so we take a fresh write ts per allocation: two commits must not
738    /// share a timestamp.
739    async fn allocate_replica_ids_for_creates(
740        &mut self,
741        decisions: &[Decision],
742    ) -> Option<Vec<ReplicaId>> {
743        let mut replica_ids = Vec::new();
744        for decision in decisions {
745            let Decision::CreateReplica { cluster_id, .. } = decision else {
746                continue;
747            };
748            let id_ts = self.get_catalog_write_ts().await;
749            let result = self
750                .catalog()
751                .allocate_replica_ids(*cluster_id, 1, id_ts)
752                .await;
753            match result {
754                Ok(ids) => {
755                    replica_ids.push(ids.into_iter().next().expect("allocated one replica id"))
756                }
757                Err(err) => {
758                    warn!(%cluster_id, "cluster controller could not allocate replica id: {err}");
759                    return None;
760                }
761            }
762        }
763        Some(replica_ids)
764    }
765
766    /// Turn a decision batch into the catalog mutation ops to transact, consuming
767    /// the `replica_ids` pre-allocated for the creates (one per `CreateReplica`,
768    /// in order). Returns `None` if a target cluster has vanished or gone
769    /// unmanaged, which makes the batch incoherent and rejects it.
770    fn build_mutation_ops(
771        &self,
772        decisions: Vec<Decision>,
773        replica_ids: Vec<ReplicaId>,
774    ) -> Option<Vec<Op>> {
775        let mut replica_ids = replica_ids.into_iter();
776        let mut mutations = Vec::new();
777        let mut drops = Vec::new();
778        for decision in decisions {
779            match decision {
780                Decision::UpdateClusterState {
781                    cluster_id, write, ..
782                } => match self.build_update_cluster_config_op(cluster_id, &write) {
783                    Some(op) => mutations.push(op),
784                    // The cluster vanished. The batch is no longer coherent.
785                    None => return None,
786                },
787                Decision::CreateReplica {
788                    cluster_id,
789                    name,
790                    shape,
791                    reason,
792                    ..
793                } => {
794                    let replica_id = replica_ids.next().expect("one pre-allocated id per create");
795                    let reason = audit_reason_for_create(reason);
796                    match self.build_create_replica_op(cluster_id, replica_id, name, &shape, reason)
797                    {
798                        Ok(Some(op)) => mutations.push(op),
799                        Ok(None) => return None,
800                        Err(err) => {
801                            warn!(%cluster_id, "cluster controller could not build replica create: {err}");
802                            return None;
803                        }
804                    }
805                }
806                Decision::DropReplica {
807                    cluster_id,
808                    replica_id,
809                    ..
810                } => {
811                    // The replica may have vanished since the decisions were
812                    // derived (a user DDL landed between the tick's read and
813                    // this apply). The in-transaction witness check would
814                    // reject such a stale batch, but resource-limit validation
815                    // runs before the transaction and panics on a missing
816                    // replica, so reject the batch here instead.
817                    if self
818                        .catalog()
819                        .try_get_cluster_replica(cluster_id, replica_id)
820                        .is_none()
821                    {
822                        return None;
823                    }
824                    drops.push(DropObjectInfo::ClusterReplica((
825                        cluster_id,
826                        replica_id,
827                        ReplicaCreateDropReason::Retired,
828                    )));
829                }
830            }
831        }
832        if !drops.is_empty() {
833            mutations.push(Op::DropObjects(drops));
834        }
835        Some(mutations)
836    }
837
838    /// Prepend the per-cluster compare-and-append `checks` to `mutations` and
839    /// transact them together.
840    ///
841    /// The checks run inside the transaction, before any mutation, so they cannot
842    /// be separated from the commit they guard. A cluster whose durable state has
843    /// diverged from what the decisions were derived from (e.g. a user `ALTER`
844    /// landed mid-tick) aborts the whole batch, so a stale create or drop can
845    /// never reshape the replica set against the config the `ALTER` has since
846    /// established (in particular, a stale drop cannot retire a replica the
847    /// `ALTER` has just made desired). On rejection nothing is applied.
848    async fn commit_with_checks(
849        &mut self,
850        checks: Vec<(ClusterId, ExpectedClusterState)>,
851        mutations: Vec<Op>,
852    ) -> ApplyOutcome {
853        let mut ops: Vec<Op> = checks
854            .into_iter()
855            .map(|(cluster_id, expected)| Op::CheckClusterState {
856                cluster_id,
857                expected,
858            })
859            .collect();
860        ops.extend(mutations);
861
862        match self.catalog_transact(None, ops).await {
863            Ok(()) => ApplyOutcome::Applied,
864            Err(AdapterError::ClusterStateChanged { .. }) => {
865                // A concurrent `ALTER` moved a cluster's durable state out from
866                // under the decisions. Expected, so the controller recomputes
867                // next tick.
868                ApplyOutcome::Rejected
869            }
870            Err(AdapterError::ReadOnly) => {
871                // The controller is quiesced while read-only (see
872                // `handle_cluster_controller_request`), so this is normally
873                // unreachable; if reached it's expected and not actionable, not
874                // a failure to surface.
875                debug!("cluster controller apply skipped in read-only mode");
876                ApplyOutcome::Rejected
877            }
878            Err(AdapterError::ResourceExhaustion { .. }) => {
879                // The batch cannot fit the resource budget. Report the fact and
880                // leave the reaction (what, if anything, to shed) to the kernel.
881                debug!("cluster controller apply exceeded the resource budget");
882                ApplyOutcome::ResourceExhausted
883            }
884            Err(err) => {
885                warn!("cluster controller apply failed: {err}");
886                ApplyOutcome::Rejected
887            }
888        }
889    }
890
891    /// Build an [`Op::UpdateClusterConfig`] that applies `write`'s deltas to the
892    /// cluster's current in-memory config, or `None` if the cluster is gone or
893    /// unmanaged. The write was guard-checked against the same state, so this is
894    /// the realized cut-over / record write.
895    fn build_update_cluster_config_op(
896        &self,
897        cluster_id: ClusterId,
898        write: &StateWrite,
899    ) -> Option<Op> {
900        let cluster = self.catalog().try_get_cluster(cluster_id)?;
901        let mut config = cluster.config.clone();
902        let ClusterConfig {
903            variant: ClusterVariant::Managed(managed),
904            ..
905        } = &mut config
906        else {
907            return None;
908        };
909        // Exhaustive destructure of the source (no `..`): a field added to
910        // `StateWrite` is a compile error here until it's overlaid onto the
911        // managed config. We cannot destructure `managed` itself. It carries
912        // fields the controller does not model (`workload_class`,
913        // `optimizer_feature_overrides`) that this overlay must leave untouched.
914        let StateWrite {
915            new_size,
916            new_replication_factor,
917            new_availability_zones,
918            new_logging,
919            new_arrangement_compression,
920            reconfiguration,
921            burst,
922        } = write;
923        if let Some(size) = new_size {
924            managed.size = size.clone();
925        }
926        if let Some(rf) = new_replication_factor {
927            managed.replication_factor = *rf;
928        }
929        if let Some(azs) = new_availability_zones {
930            managed.availability_zones = azs.clone();
931        }
932        if let Some(logging) = new_logging {
933            managed.logging = logging.clone();
934        }
935        if let Some(arrangement_compression) = new_arrangement_compression {
936            managed.arrangement_compression = *arrangement_compression;
937        }
938        if let Some(reconfiguration) = reconfiguration {
939            managed.reconfiguration = reconfiguration.record.as_ref().map(memory_reconfiguration);
940        }
941        if let Some(burst) = burst {
942            managed.burst = burst.record.as_ref().map(memory_burst);
943        }
944        // The audit intents travel with the write, declared by the strategy at
945        // the decision point. We pass them through untouched so the events are
946        // emitted in the same catalog transaction as the state they describe.
947        let reconfiguration_audit = write.reconfiguration.as_ref().and_then(|w| w.audit);
948        let burst_audit = write.burst.as_ref().and_then(|w| w.audit);
949        Some(Op::UpdateClusterConfig {
950            id: cluster_id,
951            name: cluster.name.clone(),
952            config,
953            reconfiguration_audit,
954            burst_audit,
955        })
956    }
957
958    /// Build an [`Op::CreateClusterReplica`] for a desired replica `shape` on
959    /// `cluster_id` with the pre-allocated `replica_id`, attributed to `reason`.
960    /// Returns `Ok(None)` if the cluster is gone or unmanaged.
961    fn build_create_replica_op(
962        &self,
963        cluster_id: ClusterId,
964        replica_id: ReplicaId,
965        name: String,
966        shape: &ReplicaShape,
967        reason: ReplicaCreateDropReason,
968    ) -> Result<Option<Op>, mz_catalog::memory::error::Error> {
969        let Some(cluster) = self.catalog().try_get_cluster(cluster_id) else {
970            return Ok(None);
971        };
972        if !cluster.is_managed() {
973            return Ok(None);
974        }
975        let owner_id = cluster.owner_id;
976
977        let location = mz_catalog::durable::ReplicaLocation::Managed {
978            // Concretized from the cluster config below; left empty here.
979            availability_zones: Vec::new(),
980            billed_as: None,
981            internal: false,
982            size: shape.size.clone(),
983            pending: false,
984        };
985        let azs: Option<&[String]> = if shape.availability_zones.0.is_empty() {
986            None
987        } else {
988            Some(&shape.availability_zones.0)
989        };
990        let location = self.catalog().concretize_replica_location(
991            location,
992            &self
993                .catalog()
994                .get_role_allowed_cluster_sizes(&Some(owner_id)),
995            azs,
996            false,
997        )?;
998
999        let config = mz_controller::clusters::ReplicaConfig {
1000            location,
1001            compute: ComputeReplicaConfig {
1002                logging: shape.logging.clone(),
1003                arrangement_compression: shape.arrangement_compression,
1004            },
1005        };
1006
1007        Ok(Some(Op::CreateClusterReplica {
1008            cluster_id,
1009            replica_id,
1010            name,
1011            config,
1012            owner_id,
1013            reason,
1014        }))
1015    }
1016}
1017
1018/// Map a create decision's [`CreateReason`] to the audit reason carried on the
1019/// create event. `Baseline` audits [`ReplicaCreateDropReason::Manual`], the tag
1020/// for replicas the user's own cluster config calls for. The match is
1021/// exhaustive, so a new `CreateReason` variant is a compile error here instead
1022/// of a silent `Manual`.
1023///
1024/// Drops never come through here: a drop happens exactly when no strategy
1025/// desires the replica, so it carries no attribution and is uniformly audited
1026/// [`ReplicaCreateDropReason::Retired`].
1027fn audit_reason_for_create(reason: CreateReason) -> ReplicaCreateDropReason {
1028    match reason {
1029        CreateReason::Baseline => ReplicaCreateDropReason::Manual,
1030        CreateReason::GracefulReconfiguration => ReplicaCreateDropReason::GracefulReconfiguration,
1031        CreateReason::HydrationBurst => ReplicaCreateDropReason::HydrationBurst,
1032        CreateReason::OnRefresh(decision) => ReplicaCreateDropReason::OnRefresh(decision),
1033    }
1034}
1035
1036/// Map an in-memory replica config to a [`ReplicaShape`], or `None` for an
1037/// unmanaged replica (which the controller does not own).
1038fn replica_shape(config: &mz_controller::clusters::ReplicaConfig) -> Option<ReplicaShape> {
1039    use mz_controller::clusters::ReplicaLocation;
1040    let ReplicaLocation::Managed(managed) = &config.location else {
1041        return None;
1042    };
1043    Some(ReplicaShape {
1044        size: managed.size.clone(),
1045        availability_zones: AvailabilityZones(managed.availability_zones.clone()),
1046        logging: config.compute.logging.clone(),
1047        arrangement_compression: config.compute.arrangement_compression,
1048    })
1049}
1050
1051fn on_timeout_from_controller(action: OnTimeout) -> mz_sql::plan::OnTimeoutAction {
1052    match action {
1053        OnTimeout::Commit => mz_sql::plan::OnTimeoutAction::Commit,
1054        OnTimeout::Rollback => mz_sql::plan::OnTimeoutAction::Rollback,
1055    }
1056}
1057
1058fn memory_reconfiguration(
1059    record: &ReconfigurationRecord,
1060) -> mz_catalog::memory::objects::ReconfigurationState {
1061    // Destructure the source (no `..`): a field added to the controller type is a
1062    // compile error here until it's carried across. The target is the same.
1063    let ReconfigurationRecord {
1064        target,
1065        deadline,
1066        on_timeout,
1067        status,
1068    } = record;
1069    let ReconfigurationTarget {
1070        size,
1071        replication_factor,
1072        availability_zones,
1073        logging,
1074        arrangement_compression,
1075    } = target;
1076    mz_catalog::memory::objects::ReconfigurationState {
1077        target: mz_catalog::memory::objects::ReconfigurationTarget {
1078            size: size.clone(),
1079            replication_factor: *replication_factor,
1080            availability_zones: availability_zones.0.clone(),
1081            logging: logging.clone(),
1082            arrangement_compression: *arrangement_compression,
1083        },
1084        deadline: *deadline,
1085        on_timeout: on_timeout_from_controller(*on_timeout),
1086        status: status_from_controller(*status),
1087    }
1088}
1089
1090fn status_from_controller(
1091    status: ReconfigurationStatus,
1092) -> mz_catalog::memory::objects::ReconfigurationStatus {
1093    match status {
1094        ReconfigurationStatus::InProgress => {
1095            mz_catalog::memory::objects::ReconfigurationStatus::InProgress
1096        }
1097        ReconfigurationStatus::Finalized => {
1098            mz_catalog::memory::objects::ReconfigurationStatus::Finalized
1099        }
1100        ReconfigurationStatus::TimedOut => {
1101            mz_catalog::memory::objects::ReconfigurationStatus::TimedOut
1102        }
1103        ReconfigurationStatus::Cancelled => {
1104            mz_catalog::memory::objects::ReconfigurationStatus::Cancelled
1105        }
1106        ReconfigurationStatus::ResourceExhausted => {
1107            mz_catalog::memory::objects::ReconfigurationStatus::ResourceExhausted
1108        }
1109    }
1110}
1111
1112fn memory_burst(
1113    record: &mz_cluster_controller::ctx::BurstRecord,
1114) -> mz_catalog::memory::objects::BurstState {
1115    // Destructure the source (no `..`): a field added to the controller type is a
1116    // compile error here until it's carried across.
1117    let mz_cluster_controller::ctx::BurstRecord {
1118        burst_size,
1119        linger_duration,
1120        steady_hydrated_at,
1121    } = record;
1122    mz_catalog::memory::objects::BurstState {
1123        burst_size: burst_size.clone(),
1124        linger_duration: *linger_duration,
1125        steady_hydrated_at: *steady_hydrated_at,
1126    }
1127}
1128
1129#[cfg(test)]
1130mod tests {
1131    use super::*;
1132
1133    #[mz_ore::test]
1134    fn test_audit_reason_for_create() {
1135        use ReplicaCreateDropReason as Reason;
1136        use mz_cluster_controller::ctx::RefreshWindowDecision;
1137        use mz_repr::GlobalId;
1138
1139        // Each variant maps to its own audit reason, with the baseline
1140        // auditing `Manual`.
1141        assert!(matches!(
1142            audit_reason_for_create(CreateReason::Baseline),
1143            Reason::Manual
1144        ));
1145        assert!(matches!(
1146            audit_reason_for_create(CreateReason::GracefulReconfiguration),
1147            Reason::GracefulReconfiguration
1148        ));
1149        assert!(matches!(
1150            audit_reason_for_create(CreateReason::HydrationBurst),
1151            Reason::HydrationBurst
1152        ));
1153
1154        // The on-refresh reason carries the create's window decision through
1155        // to the audit detail intact.
1156        let decision = RefreshWindowDecision {
1157            objects_needing_refresh: vec![GlobalId::User(1)],
1158            objects_needing_compaction: vec![GlobalId::User(2)],
1159            hydration_time_estimate: Duration::from_secs(7),
1160        };
1161        match audit_reason_for_create(CreateReason::OnRefresh(decision.clone())) {
1162            Reason::OnRefresh(carried) => assert_eq!(carried, decision),
1163            other => panic!("expected an on-refresh reason, got {other:?}"),
1164        }
1165    }
1166}