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}