Skip to main content

mz_cluster_controller/
strategy.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//! The pure strategy interface and the strategy implementations.
11//!
12//! A strategy is two pure functions over `(observed cluster state, live
13//! signals, now)`:
14//!
15//! - [`Strategy::update_state`] returns the durable writes the strategy wants
16//!   (cut-overs, record writes/clears). The controller transacts these in the
17//!   tick's first phase.
18//! - [`Strategy::desired_replicas`] returns the replica slots the strategy
19//!   contributes to the cluster's desired set. The controller unions every
20//!   strategy's contribution in the tick's second phase.
21//!
22//! Both are pure: same inputs, same output, no I/O. The controller is the sole
23//! mutator. Strategies never touch the [`ClusterControllerCtx`] directly. They
24//! declare the live signals they need via [`Strategy::signal_request`] and the
25//! controller fetches those before evaluating them.
26//!
27//! [`ClusterControllerCtx`]: crate::ctx::ClusterControllerCtx
28
29use std::collections::BTreeSet;
30use std::time::Duration;
31
32use mz_controller_types::ReplicaId;
33use mz_repr::{Timestamp, TimestampManipulation};
34
35use crate::ctx::{
36    AvailabilityZones, BurstAudit, BurstFinishCause, BurstRecord, BurstWrite, ClusterSchedule,
37    ClusterState, CreateReason, OnTimeout, ReconfigurationAudit, ReconfigurationRecord,
38    ReconfigurationStatus, ReconfigurationWrite, RefreshWindowDecision, RefreshWindowInputs,
39    ReplicaShape, StateWrite,
40};
41
42/// A replica slot a strategy desires this tick. The reconcile kernel unions
43/// slots across strategies and matches them by [`ReplicaShape`] against the
44/// actual replica set.
45#[derive(Clone, Debug)]
46pub struct DesiredReplica {
47    pub shape: ReplicaShape,
48    /// Why the strategy desires the slot. Carried through the kernel onto the
49    /// create decision a slot may produce (per shape, the highest-precedence
50    /// reason among the contributing slots wins).
51    pub reason: CreateReason,
52}
53
54/// One cluster-autoscaling strategy: a pair of pure functions the controller
55/// runs each tick. See the module docs.
56///
57/// `Send + Sync` so the controller (which holds a set of boxed strategies) can
58/// run on its own task.
59pub trait Strategy: Send + Sync {
60    /// The live signals this strategy needs to evaluate `state` this tick,
61    /// declared as a pure function of the durable state and the tick's config
62    /// signals. The kernel unions the requests across strategies, fetches them
63    /// through the ctx, and passes the result to [`Strategy::update_state`] and
64    /// [`Strategy::desired_replicas`]. The default requests nothing, which suits
65    /// a strategy that works off durable state alone (like the baseline).
66    fn signal_request(&self, _state: &ClusterState, _config: &ConfigSignals) -> SignalRequest {
67        SignalRequest::default()
68    }
69
70    /// The durable writes this strategy wants for `state` at time `now`. The
71    /// default is no write, which suits a strategy that only ever contributes
72    /// replicas (like the baseline). An empty [`StateWrite`] means "write
73    /// nothing": the kernel drops it without emitting a decision.
74    fn update_state(
75        &self,
76        _state: &ClusterState,
77        _signals: &LiveSignals,
78        _config: &ConfigSignals,
79        _now: Timestamp,
80    ) -> StateWrite {
81        StateWrite::default()
82    }
83
84    /// The replica slots this strategy contributes to `state`'s desired set at
85    /// time `now`.
86    fn desired_replicas(
87        &self,
88        state: &ClusterState,
89        signals: &LiveSignals,
90        config: &ConfigSignals,
91        now: Timestamp,
92    ) -> Vec<DesiredReplica>;
93}
94
95/// The live signals a strategy asks the kernel to fetch before evaluating a
96/// cluster, declared through [`Strategy::signal_request`].
97///
98/// Live signals are observations (hydration and the like) that are not durable
99/// state, so they never participate in the compare-and-append witness. Keeping
100/// them out of [`ClusterState`] keeps that type exactly the witness material
101/// plus the observed replica set.
102#[derive(Clone, Debug, Default, PartialEq, Eq)]
103pub struct SignalRequest {
104    /// Probe which of the cluster's replicas report all collections hydrated.
105    pub hydration: bool,
106    /// Probe which of the cluster's replicas are ready to be cut over to:
107    /// hydrated, and within the configured lag of the *reference* replicas
108    /// named here, the ones the cut-over will drop. `None` does not probe.
109    ///
110    /// Surviving replicas must not raise the reference, since cut-over cannot
111    /// lose their progress.
112    pub readiness: Option<BTreeSet<ReplicaId>>,
113    /// Check whether the cluster has at least one hydratable object bound to
114    /// it. See `ClusterControllerCtx::has_hydratable_objects` for what counts.
115    pub hydratable_objects: bool,
116    /// Pull the refresh-window inputs (bound REFRESH MV frontiers, schedules,
117    /// the current read timestamp).
118    pub refresh_window: bool,
119}
120
121impl SignalRequest {
122    /// The union of two requests: a signal is fetched if any strategy asks.
123    pub fn union(self, other: SignalRequest) -> SignalRequest {
124        // Exhaustive destructure (no `..`): a signal added to the request is a
125        // compile error here until its union is spelled out.
126        let SignalRequest {
127            hydration,
128            readiness,
129            hydratable_objects,
130            refresh_window,
131        } = other;
132        SignalRequest {
133            hydration: self.hydration || hydration,
134            readiness: match (self.readiness, readiness) {
135                (None, other) | (other, None) => other,
136                (Some(mut mine), Some(theirs)) => {
137                    mine.extend(theirs);
138                    Some(mine)
139                }
140            },
141            hydratable_objects: self.hydratable_objects || hydratable_objects,
142            refresh_window: self.refresh_window || refresh_window,
143        }
144    }
145}
146
147/// Environment-wide configuration the strategies consult, latched by the kernel
148/// once per tick from the controller's dyncfgs so every strategy decides against
149/// one consistent config. Not durable cluster state, so never witness material.
150#[derive(Clone, Debug, Default, PartialEq, Eq)]
151pub struct ConfigSignals {
152    /// Whether the hydration-burst strategy is enabled environment-wide (the
153    /// break-glass flag).
154    pub burst_enabled: bool,
155    /// The system-default burst linger duration, written into a new `burst`
156    /// record when the policy's `linger_duration` is omitted.
157    pub default_burst_linger: Duration,
158}
159
160/// The fulfilled live signals for one cluster, fetched by the kernel per the
161/// unioned [`SignalRequest`] and passed alongside [`ClusterState`].
162///
163/// A signal nobody requested is left at its empty default, so a strategy must
164/// only read what it declared in [`Strategy::signal_request`].
165#[derive(Clone, Debug, Default, PartialEq, Eq)]
166pub struct LiveSignals {
167    /// The replicas observed this tick to be online and to have *all* current
168    /// collections on the cluster hydrated.
169    pub hydrated_replicas: BTreeSet<ReplicaId>,
170    /// The replicas observed this tick to be ready to cut over to: hydrated, and
171    /// within the configured lag of the reference replicas the request named,
172    /// the ones the cut-over will drop. Empty when not requested. The hydration
173    /// signal is populated independently, only when requested.
174    pub ready_replicas: BTreeSet<ReplicaId>,
175    /// Whether the cluster has at least one hydratable object. `false` when not
176    /// requested.
177    pub has_hydratable_objects: bool,
178    /// The refresh-window inputs. `None` when not requested, or when the
179    /// cluster was gone, unmanaged, or no longer scheduled `ON REFRESH` when
180    /// the ctx pulled (see [`ClusterControllerCtx::refresh_window_inputs`]).
181    ///
182    /// [`ClusterControllerCtx::refresh_window_inputs`]:
183    ///     crate::ctx::ClusterControllerCtx::refresh_window_inputs
184    pub refresh_window: Option<RefreshWindowInputs>,
185}
186
187/// The implicit baseline strategy, always present.
188///
189/// Desires `replication_factor` replicas at the cluster's realized shape
190/// (`cluster.size` plus its AZ pool, logging, and arrangement compression). It
191/// holds the steady-state set so that policy strategies normally only add to
192/// it. With only the baseline engaged, the desired set equals the realized set,
193/// so a steady-state managed cluster reconciles to no decisions.
194///
195/// The baseline holds the set only for MANUAL clusters. On a scheduled cluster
196/// the controller (not the user's `replication_factor`) owns the replica set,
197/// so the baseline desires nothing there and the on-refresh strategy is the sole
198/// contributor. (The on-refresh strategy also normalizes a scheduled cluster's
199/// `replication_factor` to `0` via `update_state`, so the two views agree after
200/// the first tick regardless.)
201///
202/// The one case where the baseline steps aside is a forced cut-over, see
203/// `forced_cutover_pending`.
204#[derive(Clone, Copy, Debug, Default)]
205pub struct BaselineStrategy;
206
207/// Whether a forced cut-over is imminent: an in-progress reconfiguration is
208/// past its deadline under `ON TIMEOUT COMMIT`, so the next cut-over commits
209/// the target whether or not it hydrated.
210///
211/// In that window the baseline yields its realized-shape replicas. Overlapping
212/// the two sets only buys availability while the target hydrates, and a forced
213/// cut-over has given up on hydration. Yielding turns the reshape into one
214/// transaction that retires the realized replicas and creates the target's, so
215/// it has to fit the larger of the two shapes rather than their sum. That is
216/// what lets a resize succeed on a budget that has no room for overlap, and it
217/// is the only way to shrink a cluster that is already near its limit.
218///
219/// If that single transaction still does not fit, it is rejected whole and the
220/// record is left in progress for `ClusterController::shed_decision` to shed,
221/// so an unaffordable target stays observable rather than half-applied.
222fn forced_cutover_pending(state: &ClusterState, now: Timestamp) -> bool {
223    state.reconfiguration.as_ref().is_some_and(|record| {
224        record.is_in_progress()
225            && now >= record.deadline
226            && matches!(record.on_timeout, OnTimeout::Commit)
227    })
228}
229
230impl Strategy for BaselineStrategy {
231    fn desired_replicas(
232        &self,
233        state: &ClusterState,
234        _signals: &LiveSignals,
235        _config: &ConfigSignals,
236        now: Timestamp,
237    ) -> Vec<DesiredReplica> {
238        if !matches!(state.schedule, ClusterSchedule::Manual) {
239            return Vec::new();
240        }
241        if forced_cutover_pending(state, now) {
242            return Vec::new();
243        }
244        let shape = state.realized_shape();
245        (0..state.replication_factor)
246            .map(|_| DesiredReplica {
247                shape: shape.clone(),
248                reason: CreateReason::Baseline,
249            })
250            .collect()
251    }
252}
253
254/// The graceful (zero-downtime) reconfiguration strategy.
255///
256/// Engaged whenever the durable `reconfiguration` record is in progress. It
257/// desires `target.replication_factor` replicas at the target shape in addition
258/// to the baseline's realized-shape replicas, so both sets serve while the new
259/// one hydrates and catches up. Once rf-many target replicas are present and ready,
260/// `update_state` cuts over: the realized config advances to the target, the
261/// record is marked finalized, and the old replicas fall out of the union and
262/// are dropped. Success takes precedence over the deadline. On a timeout,
263/// `Commit` cuts over once the complete target set exists without waiting for
264/// readiness, and the baseline stops contributing in that window so the two
265/// sets swap in one transaction rather than overlapping (see
266/// `forced_cutover_pending`). `Rollback` (the default) marks the record timed
267/// out without touching the realized config and stops desiring the target
268/// replicas, reverting to the pre-reconfiguration set.
269///
270/// Both functions are pure over the observed [`ClusterState`] and the fetched
271/// [`LiveSignals`]. Readiness is requested via [`Strategy::signal_request`]
272/// exactly while an in-progress reconfiguration is present.
273#[derive(Clone, Copy, Debug, Default)]
274pub struct GracefulReconfigurationStrategy;
275
276impl GracefulReconfigurationStrategy {
277    /// Whether the cut-over precondition holds: at least
278    /// `target.replication_factor` replicas of the target shape report ready.
279    ///
280    /// Requiring rf-many ready replicas (not just one) preserves the
281    /// high-availability guarantee of `replication_factor > 1` across the
282    /// cut-over. Extra target-shape replicas beyond the rf do not block: the
283    /// post-cut-over reconcile retires them anyway, so waiting for them to
284    /// become ready would only delay the cut-over.
285    fn target_ready(
286        &self,
287        state: &ClusterState,
288        signals: &LiveSignals,
289        record: &ReconfigurationRecord,
290    ) -> bool {
291        let target_shape = record.target.shape();
292        let ready_target_replicas = state
293            .replicas
294            .iter()
295            .filter(|r| r.owned_shape().is_some_and(|s| s.matches(&target_shape)))
296            .filter(|r| signals.ready_replicas.contains(&r.replica_id))
297            .count();
298        let target_rf = usize::try_from(record.target.replication_factor).unwrap_or(usize::MAX);
299        ready_target_replicas >= target_rf
300    }
301
302    /// Whether the complete target set exists, without requiring hydration.
303    fn target_materialized(&self, state: &ClusterState, record: &ReconfigurationRecord) -> bool {
304        let target_shape = record.target.shape();
305        let target_replicas = state
306            .replicas
307            .iter()
308            .filter(|r| r.owned_shape().is_some_and(|s| s.matches(&target_shape)))
309            .count();
310        let target_rf = usize::try_from(record.target.replication_factor).unwrap_or(usize::MAX);
311        target_replicas >= target_rf
312    }
313}
314
315impl Strategy for GracefulReconfigurationStrategy {
316    fn signal_request(&self, state: &ClusterState, _config: &ConfigSignals) -> SignalRequest {
317        let in_progress = state
318            .reconfiguration
319            .as_ref()
320            .is_some_and(|record| record.is_in_progress());
321        if !in_progress {
322            return SignalRequest::default();
323        }
324        // Only the realized-shape replicas are retired by this strategy.
325        let realized = state.realized_shape();
326        let reference = state
327            .replicas
328            .iter()
329            .filter(|r| r.owned_shape().is_some_and(|s| s.matches(&realized)))
330            .map(|r| r.replica_id)
331            .collect();
332        SignalRequest {
333            readiness: Some(reference),
334            ..Default::default()
335        }
336    }
337
338    fn update_state(
339        &self,
340        state: &ClusterState,
341        signals: &LiveSignals,
342        _config: &ConfigSignals,
343        now: Timestamp,
344    ) -> StateWrite {
345        let Some(record) = &state.reconfiguration else {
346            return StateWrite::default();
347        };
348        if !record.is_in_progress() {
349            return StateWrite::default();
350        }
351
352        // Cut over by advancing the realized config to the target and marking
353        // the record finalized on either of two conditions:
354        //   1. rf-many target replicas are present and ready (success, which
355        //      takes precedence over the deadline regardless of `on_timeout`), or
356        //   2. the deadline has been reached, `on_timeout` is `Commit`, and the
357        //      complete target set exists (cut over without waiting for readiness).
358        //
359        // NOTE: the deadline is reached at `now >= deadline`, not `now > deadline`.
360        // An `ON TIMEOUT COMMIT` with a zero timeout writes `deadline = now` to
361        // request an immediate cut-over. With a strict `>`, a first tick landing at
362        // exactly that timestamp would miss the deadline, so phase 2 would provision
363        // the overlap target replicas and only a later tick would cut over. `>=`
364        // fires the deadline the instant it is reached, so the zero-timeout cut-over
365        // happens on the first tick, before any overlap replica is desired.
366        // We require the target set to exist before a forced cut-over so its
367        // concrete create transaction can enforce resource limits. Otherwise a
368        // zero-timeout commit could finalize first, fail to create the new
369        // baseline, and leave no in-progress strategy for the controller to shed.
370        // The baseline yields while we wait (see `forced_cutover_pending`), so
371        // that create arrives in the same transaction that retires the realized
372        // replicas and does not have to fit alongside them.
373        let ready = self.target_ready(state, signals, record);
374        let deadline_reached = now >= record.deadline;
375        let commit_on_timeout = deadline_reached && matches!(record.on_timeout, OnTimeout::Commit);
376        let target_materialized = self.target_materialized(state, record);
377        if ready || (commit_on_timeout && target_materialized) {
378            return StateWrite {
379                new_size: Some(record.target.size.clone()),
380                new_replication_factor: Some(record.target.replication_factor),
381                new_availability_zones: Some(record.target.availability_zones.0.clone()),
382                new_logging: Some(record.target.logging.clone()),
383                new_arrangement_compression: Some(record.target.arrangement_compression),
384                reconfiguration: Some(ReconfigurationWrite {
385                    record: Some(ReconfigurationRecord {
386                        status: ReconfigurationStatus::Finalized,
387                        ..record.clone()
388                    }),
389                    // A cut-over that only happens because the deadline passed
390                    // under `Commit` is forced: the target is not ready.
391                    // Declared here because only this decision point knows.
392                    // The durable status reads `Finalized` either way.
393                    audit: Some(ReconfigurationAudit::Finalized { forced: !ready }),
394                }),
395                ..Default::default()
396            };
397        }
398
399        // Past the deadline not ready under `Rollback`: abandon the
400        // reconfiguration while leaving the realized config untouched. The
401        // terminal status is the durable transition the audit event records. With
402        // the record no longer in progress the strategy stops contributing the
403        // target set, so the baseline alone shapes the cluster.
404        if deadline_reached && matches!(record.on_timeout, OnTimeout::Rollback) {
405            return StateWrite {
406                reconfiguration: Some(ReconfigurationWrite {
407                    record: Some(ReconfigurationRecord {
408                        status: ReconfigurationStatus::TimedOut,
409                        ..record.clone()
410                    }),
411                    audit: Some(ReconfigurationAudit::TimedOut),
412                }),
413                ..Default::default()
414            };
415        }
416
417        // Before the deadline: keep waiting.
418        StateWrite::default()
419    }
420
421    fn desired_replicas(
422        &self,
423        state: &ClusterState,
424        signals: &LiveSignals,
425        _config: &ConfigSignals,
426        now: Timestamp,
427    ) -> Vec<DesiredReplica> {
428        let Some(record) = &state.reconfiguration else {
429            return Vec::new();
430        };
431        if !record.is_in_progress() {
432            return Vec::new();
433        }
434
435        // Past the deadline with the target not ready under `Rollback`: stop
436        // contributing the target replicas. `update_state` marks the record
437        // timed out in this same tick's first phase, so this usually never fires
438        // against a re-read state. It matters when the deadline crosses between
439        // the two phases' `ctx.now()` reads within one tick: phase 1 saw the
440        // deadline unreached and wrote nothing, phase 2 sees it reached here and
441        // already stops desiring the target, keeping the rollback's replica
442        // drops prompt rather than waiting a tick for the status write.
443        // Everything else (before the deadline, awaiting a success cut-over
444        // past it, or a `Commit` cut-over `update_state` performs this tick)
445        // keeps desiring the target set.
446        // `now >= deadline` matches `update_state`'s boundary, so a zero-timeout
447        // rollback stops desiring the target on the same tick it marks the
448        // record timed out.
449        let timed_out = now >= record.deadline && !self.target_ready(state, signals, record);
450        if timed_out && matches!(record.on_timeout, OnTimeout::Rollback) {
451            return Vec::new();
452        }
453
454        let shape = record.target.shape();
455        (0..record.target.replication_factor)
456            .map(|_| DesiredReplica {
457                shape: shape.clone(),
458                reason: CreateReason::GracefulReconfiguration,
459            })
460            .collect()
461    }
462}
463
464/// The `ON REFRESH` scheduling strategy.
465///
466/// Engaged for clusters with a non-MANUAL [`ClusterSchedule`]. It contributes one
467/// replica at the cluster's realized shape while the cluster is inside a refresh
468/// window, and nothing otherwise. The window decision keys on the bound REFRESH
469/// materialized views' write frontiers, their refresh schedules, the configured
470/// hydration-time estimate, and the current read timestamp, all carried in
471/// [`RefreshWindowInputs`].
472///
473/// The controller (not the user's `replication_factor`) owns a scheduled
474/// cluster's replica set, so [`Strategy::update_state`] normalizes the realized
475/// `replication_factor` to `0`. This is self-healing (no migration needed to
476/// enable the controller) and makes `mz_clusters.replication_factor` read `0` for
477/// a scheduled cluster, with `mz_cluster_replicas` authoritative for what is
478/// actually running.
479///
480/// NB: the decision is re-derived purely from the live signals each tick, with
481/// no cross-tick latch. We pull a complete decision from durable and storage
482/// state on every tick, so the first tick after a restart already decides from
483/// the same inputs as a steady tick.
484#[derive(Clone, Copy, Debug, Default)]
485pub struct OnRefreshStrategy;
486
487impl OnRefreshStrategy {
488    /// The window decision for the cluster: which bound REFRESH MVs either still
489    /// need a refresh (their write frontier has not advanced past the read
490    /// timestamp adjusted by the hydration-time estimate) or are estimated to
491    /// still need Persist compaction after their last refresh. The cluster
492    /// should be On iff either list is non-empty
493    /// ([`RefreshWindowDecision::window_open`]), so an open window always names
494    /// the MVs that explain it.
495    ///
496    /// `hydration_time_estimate` comes from the schedule; the remaining signals
497    /// come from `inputs`. With no bound REFRESH MVs both lists are empty and
498    /// the cluster is Off.
499    fn window_decision(
500        &self,
501        hydration_time_estimate: std::time::Duration,
502        inputs: &RefreshWindowInputs,
503    ) -> RefreshWindowDecision {
504        // 1. Needs refresh: write_frontier < read_ts + hydration_time_estimate.
505        // The cluster is turned on `hydration_time_estimate` ahead of a refresh
506        // so it can rehydrate before the refresh time.
507        let read_ts_adjusted = inputs
508            .read_ts
509            .step_forward_by(&duration_to_ts(hydration_time_estimate));
510        let objects_needing_refresh = inputs
511            .refresh_mvs
512            .iter()
513            .filter(|mv| mv.write_frontier.less_than(&read_ts_adjusted))
514            .map(|mv| mv.id)
515            .collect();
516
517        // 2. Needs compaction: prev_refresh + compaction_estimate > read_ts. We
518        // keep the cluster on for a while after a refresh so Persist can compact.
519        let compaction_estimate = duration_to_ts(inputs.compaction_estimate);
520        let objects_needing_compaction = inputs
521            .refresh_mvs
522            .iter()
523            .filter(|mv| {
524                // `prev_refresh` is None in two cases, both meaning "schedule no
525                // compaction time now": no refresh has happened yet (no frontier to
526                // round down and no past `AT`), or a `REFRESH EVERY` MV with an empty
527                // write frontier (we have no wall-clock handle on its last refresh).
528                let prev_refresh = match mv.write_frontier.as_option() {
529                    Some(frontier) => frontier.round_down_minus_1(&mv.refresh_schedule),
530                    None => mv.refresh_schedule.last_refresh(),
531                };
532                prev_refresh.is_some_and(|prev| {
533                    // An estimate that overflows the timestamp space means
534                    // `prev + estimate` exceeds every possible read ts, so the
535                    // window reads as open.
536                    match prev.try_step_forward_by(&compaction_estimate) {
537                        Some(compacting_until) => compacting_until > inputs.read_ts,
538                        None => true,
539                    }
540                })
541            })
542            .map(|mv| mv.id)
543            .collect();
544
545        RefreshWindowDecision {
546            objects_needing_refresh,
547            objects_needing_compaction,
548            hydration_time_estimate,
549        }
550    }
551}
552
553impl Strategy for OnRefreshStrategy {
554    fn signal_request(&self, state: &ClusterState, _config: &ConfigSignals) -> SignalRequest {
555        SignalRequest {
556            refresh_window: !matches!(state.schedule, ClusterSchedule::Manual),
557            ..Default::default()
558        }
559    }
560
561    fn update_state(
562        &self,
563        state: &ClusterState,
564        _signals: &LiveSignals,
565        _config: &ConfigSignals,
566        _now: Timestamp,
567    ) -> StateWrite {
568        // The controller owns a scheduled cluster's replica set, so hold the
569        // realized `replication_factor` at `0`. A stale non-zero value (e.g.
570        // carried over from a cluster that was just given a schedule) would
571        // otherwise have the implicit baseline desire a replica the on-refresh
572        // strategy does not, a flap.
573        // Only write when it is actually non-zero, to keep steady ticks no-ops.
574        if matches!(state.schedule, ClusterSchedule::Manual) || state.replication_factor == 0 {
575            return StateWrite::default();
576        }
577        // While a reconfiguration record is in progress, the graceful strategy
578        // owns `new_replication_factor` (its cut-over sets it from the record's
579        // target), so skip the normalization to keep the field single-writer
580        // within a tick. The sequencer never writes a record for a scheduled
581        // cluster, so this state is reachable only for a record written before
582        // the cluster acquired its schedule (pre-upgrade catalog state). A
583        // cut-over there can briefly set a non-zero rf on the scheduled
584        // cluster. The next tick sees the record settled and normalizes it.
585        if state
586            .reconfiguration
587            .as_ref()
588            .is_some_and(|record| record.is_in_progress())
589        {
590            return StateWrite::default();
591        }
592        StateWrite {
593            new_replication_factor: Some(0),
594            ..Default::default()
595        }
596    }
597
598    fn desired_replicas(
599        &self,
600        state: &ClusterState,
601        signals: &LiveSignals,
602        _config: &ConfigSignals,
603        _now: Timestamp,
604    ) -> Vec<DesiredReplica> {
605        let ClusterSchedule::Refresh {
606            hydration_time_estimate,
607        } = state.schedule
608        else {
609            return Vec::new();
610        };
611        // The refresh-window signals are pulled for every scheduled cluster.
612        // The ctx returns `None` only when the cluster was gone, unmanaged, or
613        // no longer scheduled at pull time (a concurrent DDL moved it under the
614        // tick), so contributing nothing is the correct answer. The schedule is
615        // part of the compare-and-append witness, so a stale in-flight decision
616        // derived before such a change is rejected at apply anyway.
617        let Some(inputs) = &signals.refresh_window else {
618            return Vec::new();
619        };
620        let decision = self.window_decision(hydration_time_estimate, inputs);
621        if !decision.window_open() {
622            return Vec::new();
623        }
624        // One replica at the realized shape (`cluster.size` plus the cluster's AZ
625        // pool, logging, and arrangement compression). The window decision rides
626        // inside the reason so the create it may produce can carry the audit
627        // detail.
628        vec![DesiredReplica {
629            shape: state.realized_shape(),
630            reason: CreateReason::OnRefresh(decision),
631        }]
632    }
633}
634
635/// A millisecond [`Duration`] as a [`Timestamp`], saturating at [`Timestamp::MAX`]
636/// on overflow rather than panicking the controller on a bad input.
637///
638/// [`Duration`]: std::time::Duration
639pub fn duration_to_ts(duration: std::time::Duration) -> Timestamp {
640    Timestamp::try_from(duration).unwrap_or(Timestamp::MAX)
641}
642
643/// The hydration-burst strategy.
644///
645/// Engaged for clusters whose `AUTO SCALING STRATEGY` sets `ON HYDRATION`. While
646/// the cluster is On and there exists an object on it that no steady-state
647/// (realized-config) replica has hydrated, it runs one extra replica at the
648/// configured `HYDRATION SIZE` to accelerate hydration; the burst replica tears
649/// down a `linger_duration` after the steady set first hydrates. Zero objects
650/// make the condition vacuously unsatisfied, so a brand-new cluster never bursts
651/// before its first object lands. The burst is keyed entirely on the presence of a
652/// durable `burst` record (written/cleared by [`Strategy::update_state`]); the
653/// burst replica is an ordinary replica. The union/diff reconciler creates and
654/// drops it by shape+count with no special identity.
655///
656/// There is deliberately no TTL on the burst replica: if the steady set can never
657/// hydrate at `cluster.size`, the burst stays up indefinitely (the cluster runs
658/// permanently oversized, visible in billing and the audit log), the accepted
659/// trade for keeping the cluster serving. Burst is **not** suppressed during a
660/// reconfiguration; the two coexist.
661///
662/// Steady-replica hydration and object existence are live signals requested via
663/// [`Strategy::signal_request`] while an `ON HYDRATION` policy is active.
664#[derive(Clone, Copy, Debug, Default)]
665pub struct HydrationBurstStrategy;
666
667impl HydrationBurstStrategy {
668    /// The cluster's active `ON HYDRATION` policy, but only when burst is permitted
669    /// at all: the break-glass flag is on and the cluster is On (`rf > 0`). `None`
670    /// otherwise. No burst is warranted and any existing record is torn down.
671    fn active_policy<'a>(
672        &self,
673        state: &'a ClusterState,
674        config: &ConfigSignals,
675    ) -> Option<&'a crate::ctx::OnHydrationPolicy> {
676        if !config.burst_enabled || state.replication_factor == 0 {
677            return None;
678        }
679        state.auto_scaling_policy.as_ref()?.on_hydration.as_ref()
680    }
681
682    /// The in-flight burst record, but only while the current config still
683    /// warrants it: the policy is active ([`Self::active_policy`]) and the
684    /// record's size matches the policy's `HYDRATION SIZE`. `None` for a stale
685    /// record, which `update_state` tears down.
686    fn warranted_record<'a>(
687        &self,
688        state: &'a ClusterState,
689        config: &ConfigSignals,
690    ) -> Option<&'a BurstRecord> {
691        let record = state.burst.as_ref()?;
692        // `active_policy` already folds in `replication_factor != 0`, so the
693        // shared predicate's own check is redundant here, but passing the real
694        // value keeps this a faithful call of the one warrant definition.
695        let hydration_size = self
696            .active_policy(state, config)
697            .map(|policy| policy.hydration_size.as_str());
698        mz_adapter_types::cluster_state::burst_record_warranted(
699            &record.burst_size,
700            state.replication_factor,
701            hydration_size,
702        )
703        .then_some(record)
704    }
705
706    /// Whether at least one steady-state (realized-config) replica reports all
707    /// current objects hydrated. `false` when no steady replica reports at all
708    /// (absent, or not yet registered with the compute controller).
709    fn steady_hydrated(&self, state: &ClusterState, signals: &LiveSignals) -> bool {
710        let steady_shape = state.realized_shape();
711        state
712            .replicas
713            .iter()
714            .filter(|r| r.owned_shape().is_some_and(|s| s.matches(&steady_shape)))
715            .any(|r| signals.hydrated_replicas.contains(&r.replica_id))
716    }
717}
718
719impl Strategy for HydrationBurstStrategy {
720    fn signal_request(&self, state: &ClusterState, config: &ConfigSignals) -> SignalRequest {
721        // Hydration drives both the arm check and the linger lifecycle. Object
722        // existence only gates arming, so it is requested only record-less.
723        let active = self.active_policy(state, config).is_some();
724        SignalRequest {
725            hydration: active,
726            hydratable_objects: active && state.burst.is_none(),
727            ..Default::default()
728        }
729    }
730
731    fn update_state(
732        &self,
733        state: &ClusterState,
734        signals: &LiveSignals,
735        config: &ConfigSignals,
736        now: Timestamp,
737    ) -> StateWrite {
738        // Both teardown arms clear the record, but they declare different
739        // causes: only this decision point knows whether the burst ran its
740        // course or was cut short by a config change.
741        let clear = |cause: BurstFinishCause| StateWrite {
742            burst: Some(BurstWrite {
743                record: None,
744                audit: Some(BurstAudit::Finished { cause }),
745            }),
746            ..Default::default()
747        };
748
749        // Cleanup precedence: a burst no longer warranted tears down regardless
750        // of linger. Catalog writes retire records they invalidate themselves,
751        // so this arm mainly covers the burst dyncfg switching off, and
752        // backstops any stale record that reaches us anyway.
753        if state.burst.is_some() && self.warranted_record(state, config).is_none() {
754            return clear(BurstFinishCause::NoLongerWarranted);
755        }
756        let Some(policy) = self.active_policy(state, config) else {
757            // No record (the cleanup above handled that) and no active policy:
758            // nothing to arm.
759            return StateWrite::default();
760        };
761
762        let steady_hydrated = self.steady_hydrated(state, signals);
763
764        match &state.burst {
765            // No record: arm a burst only while some object exists that the
766            // steady set has not hydrated. Without the object gate, a brand-new
767            // cluster would burst at creation with nothing to accelerate (an
768            // absent steady replica reads as un-hydrated). The record-present
769            // arms below do not consult the gate: if all objects are dropped
770            // mid-burst, the steady set reads hydrated and the linger clears
771            // the record.
772            None => {
773                if steady_hydrated || !signals.has_hydratable_objects {
774                    StateWrite::default()
775                } else {
776                    let linger_duration = policy
777                        .linger_duration
778                        .unwrap_or(config.default_burst_linger);
779                    StateWrite {
780                        burst: Some(BurstWrite {
781                            record: Some(BurstRecord {
782                                burst_size: policy.hydration_size.clone(),
783                                linger_duration,
784                                steady_hydrated_at: None,
785                            }),
786                            audit: Some(BurstAudit::Started),
787                        }),
788                        ..Default::default()
789                    }
790                }
791            }
792            // Record present: drive the linger/teardown/re-arm lifecycle.
793            Some(record) => {
794                match (record.steady_hydrated_at, steady_hydrated) {
795                    // Steady set hydrated and the linger has elapsed: tear down.
796                    // A linger that overflows the timestamp space reads as
797                    // never-elapsed.
798                    (Some(hydrated_at), true)
799                        if now
800                            > hydrated_at
801                                .try_step_forward_by(&duration_to_ts(record.linger_duration))
802                                .unwrap_or(Timestamp::MAX) =>
803                    {
804                        clear(BurstFinishCause::LingerElapsed)
805                    }
806                    // Steady set hydrated, linger not yet elapsed: hold.
807                    (Some(_), true) => StateWrite::default(),
808                    // First observation of the steady set hydrated: stamp the
809                    // linger start. A bookkeeping rewrite, not a lifecycle
810                    // transition, so it declares no audit.
811                    (None, true) => StateWrite {
812                        burst: Some(BurstWrite {
813                            record: Some(BurstRecord {
814                                steady_hydrated_at: Some(now),
815                                ..record.clone()
816                            }),
817                            audit: None,
818                        }),
819                        ..Default::default()
820                    },
821                    // The steady set went un-hydrated again after we had stamped a
822                    // hydration time: re-arm so the linger restarts after the next
823                    // successful hydration. Also bookkeeping: the burst replica
824                    // keeps running throughout, so no lifecycle event.
825                    (Some(_), false) => StateWrite {
826                        burst: Some(BurstWrite {
827                            record: Some(BurstRecord {
828                                steady_hydrated_at: None,
829                                ..record.clone()
830                            }),
831                            audit: None,
832                        }),
833                        ..Default::default()
834                    },
835                    // Steady set still un-hydrated and never stamped: keep waiting.
836                    (None, false) => StateWrite::default(),
837                }
838            }
839        }
840    }
841
842    fn desired_replicas(
843        &self,
844        state: &ClusterState,
845        _signals: &LiveSignals,
846        _config: &ConfigSignals,
847        _now: Timestamp,
848    ) -> Vec<DesiredReplica> {
849        // A present record is never stale: catalog writes retire records they
850        // invalidate in the same transaction, and a dyncfg switch-off is
851        // handled by phase 1's cleanup (config signals are latched per tick).
852        // One replica at the burst size (only the size differs from steady).
853        let Some(record) = &state.burst else {
854            return Vec::new();
855        };
856        vec![DesiredReplica {
857            shape: ReplicaShape {
858                size: record.burst_size.clone(),
859                availability_zones: AvailabilityZones(state.availability_zones.clone()),
860                logging: state.logging.clone(),
861                arrangement_compression: state.arrangement_compression,
862            },
863            reason: CreateReason::HydrationBurst,
864        }]
865    }
866}