Skip to main content

mz_compute_client/controller/
instance.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//! A controller for a compute instance.
11
12use std::collections::{BTreeMap, BTreeSet};
13use std::fmt::Debug;
14use std::sync::{Arc, Mutex};
15use std::time::{Duration, Instant};
16
17use chrono::{DateTime, DurationRound, TimeDelta, Utc};
18use differential_dataflow::lattice::Lattice;
19use mz_build_info::BuildInfo;
20use mz_cluster_client::WallclockLagFn;
21use mz_compute_types::dataflows::{BuildDesc, DataflowDescription};
22use mz_compute_types::plan::render_plan::RenderPlan;
23use mz_compute_types::sinks::{
24    ComputeSinkConnection, ComputeSinkDesc, MaterializedViewSinkConnection,
25};
26use mz_compute_types::sources::SourceInstanceDesc;
27use mz_controller_types::dyncfgs::{
28    ENABLE_PAUSED_CLUSTER_READHOLD_DOWNGRADE, WALLCLOCK_LAG_RECORDING_INTERVAL,
29};
30use mz_dyncfg::{ConfigSet, ConfigUpdates};
31use mz_expr::RowSetFinishing;
32use mz_ore::cast::CastFrom;
33use mz_ore::channel::instrumented_unbounded_channel;
34use mz_ore::now::NowFn;
35use mz_ore::tracing::OpenTelemetryContext;
36use mz_ore::{soft_assert_or_log, soft_panic_or_log};
37use mz_persist_types::PersistLocation;
38use mz_repr::adt::timestamp::CheckedTimestamp;
39use mz_repr::refresh_schedule::RefreshSchedule;
40use mz_repr::{Datum, Diff, GlobalId, RelationDesc, Row, Timestamp};
41use mz_storage_client::controller::{IntrospectionType, WallclockLag, WallclockLagHistogramPeriod};
42use mz_storage_types::read_holds::{self, ReadHold};
43use mz_storage_types::read_policy::ReadPolicy;
44use thiserror::Error;
45use timely::PartialOrder;
46use timely::progress::frontier::MutableAntichain;
47use timely::progress::{Antichain, ChangeBatch};
48use tokio::sync::{mpsc, oneshot};
49use uuid::Uuid;
50
51use crate::controller::error::{
52    CollectionMissing, ERROR_TARGET_REPLICA_FAILED, HydrationCheckBadTarget,
53};
54use crate::controller::instance_client::PeekError;
55use crate::controller::replica::{ReplicaClient, ReplicaConfig};
56use crate::controller::{
57    CollectionReadiness, ComputeControllerResponse, IntrospectionUpdates, PeekNotification,
58    ReplicaId, StorageCollections,
59};
60use crate::logging::LogVariant;
61use crate::metrics::IntCounter;
62use crate::metrics::{InstanceMetrics, ReplicaCollectionMetrics, ReplicaMetrics, UIntGauge};
63use crate::protocol::command::{
64    ComputeCommand, ComputeParameters, InstanceConfig, Peek, PeekTarget,
65};
66use crate::protocol::history::ComputeCommandHistory;
67use crate::protocol::response::{
68    ComputeResponse, CopyToResponse, FrontiersResponse, PeekError as ProtocolPeekError,
69    PeekResponse, StatusResponse, SubscribeBatch, SubscribeResponse,
70};
71
72#[derive(Error, Debug)]
73#[error("replica exists already: {0}")]
74pub(super) struct ReplicaExists(pub ReplicaId);
75
76#[derive(Error, Debug)]
77#[error("replica does not exist: {0}")]
78pub(super) struct ReplicaMissing(pub ReplicaId);
79
80#[derive(Error, Debug)]
81pub(super) enum DataflowCreationError {
82    #[error("collection does not exist: {0}")]
83    CollectionMissing(GlobalId),
84    #[error("replica does not exist: {0}")]
85    ReplicaMissing(ReplicaId),
86    #[error("dataflow definition lacks an as_of value")]
87    MissingAsOf,
88    #[error("subscribe dataflow has an empty as_of")]
89    EmptyAsOfForSubscribe,
90    #[error("copy to dataflow has an empty as_of")]
91    EmptyAsOfForCopyTo,
92    #[error("no read hold provided for dataflow import: {0}")]
93    ReadHoldMissing(GlobalId),
94    #[error("insufficient read hold provided for dataflow import: {0}")]
95    ReadHoldInsufficient(GlobalId),
96}
97
98impl From<CollectionMissing> for DataflowCreationError {
99    fn from(error: CollectionMissing) -> Self {
100        Self::CollectionMissing(error.0)
101    }
102}
103
104#[derive(Error, Debug)]
105pub(super) enum ReadPolicyError {
106    #[error("collection does not exist: {0}")]
107    CollectionMissing(GlobalId),
108    #[error("collection is write-only: {0}")]
109    WriteOnlyCollection(GlobalId),
110}
111
112impl From<CollectionMissing> for ReadPolicyError {
113    fn from(error: CollectionMissing) -> Self {
114        Self::CollectionMissing(error.0)
115    }
116}
117
118/// A command sent to an [`Instance`] task.
119pub(super) type Command = Box<dyn FnOnce(&mut Instance) + Send>;
120
121/// A response from a replica, composed of a replica ID, the replica's current epoch, and the
122/// compute response itself.
123pub(super) type ReplicaResponse = (ReplicaId, u64, ComputeResponse);
124
125/// The state we keep for a compute instance.
126pub(super) struct Instance {
127    /// Build info for spawning replicas
128    build_info: &'static BuildInfo,
129    /// A handle providing access to storage collections.
130    storage_collections: StorageCollections,
131    /// Whether instance initialization has been completed.
132    initialized: bool,
133    /// Whether this instance is in read-only mode.
134    ///
135    /// When in read-only mode, this instance will not update persistent state, such as
136    /// wallclock lag introspection.
137    read_only: bool,
138    /// The workload class of this instance.
139    ///
140    /// This is currently only used to annotate metrics.
141    workload_class: Option<String>,
142    /// The replicas of this compute instance.
143    replicas: BTreeMap<ReplicaId, ReplicaState>,
144    /// Per-replica dyncfg overrides, merged into the `UpdateConfiguration`
145    /// command sent to each replica (and into the command-history replay used
146    /// to hydrate new replicas). Populated from the scoped feature flags
147    /// (replica-local) layer; empty by default, in which case every replica
148    /// receives the unmodified environment-wide configuration. Stores only the
149    /// values that differ from the environment-wide value, so the map is sparse.
150    replica_dyncfg_overrides: BTreeMap<ReplicaId, ConfigUpdates>,
151    /// Currently installed compute collections.
152    ///
153    /// New entries are added for all collections exported from dataflows created through
154    /// [`Instance::create_dataflow`].
155    ///
156    /// Entries are removed by [`Instance::cleanup_collections`]. See that method's documentation
157    /// about the conditions for removing collection state.
158    collections: BTreeMap<GlobalId, CollectionState>,
159    /// IDs of log sources maintained by this compute instance.
160    log_sources: BTreeMap<LogVariant, GlobalId>,
161    /// Currently outstanding peeks.
162    ///
163    /// New entries are added for all peeks initiated through [`Instance::peek`].
164    ///
165    /// The entry for a peek is only removed once all replicas have responded to the peek. This is
166    /// currently required to ensure all replicas have stopped reading from the peeked collection's
167    /// inputs before we allow them to compact. database-issues#4822 tracks changing this so we only have to wait
168    /// for the first peek response.
169    peeks: BTreeMap<Uuid, PendingPeek>,
170    /// Currently in-progress subscribes.
171    ///
172    /// New entries are added for all subscribes exported from dataflows created through
173    /// [`Instance::create_dataflow`].
174    ///
175    /// The entry for a subscribe is removed once at least one replica has reported the subscribe
176    /// to have advanced to the empty frontier or to have been dropped, implying that no further
177    /// updates will be emitted for this subscribe.
178    ///
179    /// Note that subscribes are tracked both in `collections` and `subscribes`. `collections`
180    /// keeps track of the subscribe's upper and since frontiers and ensures appropriate read holds
181    /// on the subscribe's input. `subscribes` is only used to track which updates have been
182    /// emitted, to decide if new ones should be emitted or suppressed.
183    subscribes: BTreeMap<GlobalId, ActiveSubscribe>,
184    /// Tracks all in-progress COPY TOs.
185    ///
186    /// New entries are added for all s3 oneshot sinks (corresponding to a COPY TO) exported from
187    /// dataflows created through [`Instance::create_dataflow`].
188    ///
189    /// The entry for a copy to is removed once at least one replica has finished
190    /// or the exporting collection is dropped.
191    copy_tos: BTreeSet<GlobalId>,
192    /// The command history, used when introducing new replicas or restarting existing replicas.
193    history: ComputeCommandHistory<UIntGauge>,
194    /// Receiver for commands to be executed.
195    command_rx: mpsc::UnboundedReceiver<Command>,
196    /// Sender for responses to be delivered.
197    response_tx: mpsc::UnboundedSender<ComputeControllerResponse>,
198    /// Sender for introspection updates to be recorded.
199    introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
200    /// The registry the controller uses to report metrics.
201    metrics: InstanceMetrics,
202    /// Dynamic system configuration.
203    dyncfg: Arc<ConfigSet>,
204
205    /// The persist location where we can stash large peek results.
206    peek_stash_persist_location: PersistLocation,
207
208    /// A function that produces the current wallclock time.
209    now: NowFn,
210    /// A function that computes the lag between the given time and wallclock time.
211    wallclock_lag: WallclockLagFn<Timestamp>,
212    /// The last time wallclock lag introspection was recorded.
213    wallclock_lag_last_recorded: DateTime<Utc>,
214
215    /// Sender for updates to collection read holds.
216    ///
217    /// Copies of this sender are given to [`ReadHold`]s that are created in
218    /// [`CollectionState::new`].
219    read_hold_tx: read_holds::ChangeTx,
220    /// A sender for responses from replicas.
221    replica_tx: mz_ore::channel::InstrumentedUnboundedSender<ReplicaResponse, IntCounter>,
222    /// A receiver for responses from replicas.
223    replica_rx: mz_ore::channel::InstrumentedUnboundedReceiver<ReplicaResponse, IntCounter>,
224}
225
226impl Instance {
227    /// Acquire a handle to the collection state associated with `id`.
228    fn collection(&self, id: GlobalId) -> Result<&CollectionState, CollectionMissing> {
229        self.collections.get(&id).ok_or(CollectionMissing(id))
230    }
231
232    /// Acquire a mutable handle to the collection state associated with `id`.
233    fn collection_mut(&mut self, id: GlobalId) -> Result<&mut CollectionState, CollectionMissing> {
234        self.collections.get_mut(&id).ok_or(CollectionMissing(id))
235    }
236
237    /// Acquire a handle to the collection state associated with `id`.
238    ///
239    /// # Panics
240    ///
241    /// Panics if the identified collection does not exist.
242    fn expect_collection(&self, id: GlobalId) -> &CollectionState {
243        self.collections.get(&id).expect("collection must exist")
244    }
245
246    /// Acquire a mutable handle to the collection state associated with `id`.
247    ///
248    /// # Panics
249    ///
250    /// Panics if the identified collection does not exist.
251    fn expect_collection_mut(&mut self, id: GlobalId) -> &mut CollectionState {
252        self.collections
253            .get_mut(&id)
254            .expect("collection must exist")
255    }
256
257    fn collections_iter(&self) -> impl Iterator<Item = (GlobalId, &CollectionState)> {
258        self.collections.iter().map(|(id, coll)| (*id, coll))
259    }
260
261    /// Returns an iterator over replicas that host the given collection.
262    ///
263    /// For replica-targeted collections, this returns only the target replica.
264    /// For non-targeted collections, this returns all replicas.
265    ///
266    /// Returns `Err` if the collection does not exist.
267    fn replicas_hosting(
268        &self,
269        id: GlobalId,
270    ) -> Result<impl Iterator<Item = &ReplicaState> + Clone, CollectionMissing> {
271        let target = self.collection(id)?.target_replica;
272        Ok(self
273            .replicas
274            .values()
275            .filter(move |r| target.map_or(true, |t| t == r.id)))
276    }
277
278    /// Add a collection to the instance state.
279    ///
280    /// # Panics
281    ///
282    /// Panics if a collection with the same ID exists already.
283    fn add_collection(
284        &mut self,
285        id: GlobalId,
286        as_of: Antichain<Timestamp>,
287        shared: SharedCollectionState,
288        storage_dependencies: BTreeMap<GlobalId, ReadHold>,
289        compute_dependencies: BTreeMap<GlobalId, ReadHold>,
290        replica_input_read_holds: Vec<ReadHold>,
291        write_only: bool,
292        storage_sink: bool,
293        initial_as_of: Option<Antichain<Timestamp>>,
294        refresh_schedule: Option<RefreshSchedule>,
295        target_replica: Option<ReplicaId>,
296    ) {
297        // Add global collection state.
298        let dependency_ids: Vec<GlobalId> = compute_dependencies
299            .keys()
300            .chain(storage_dependencies.keys())
301            .copied()
302            .collect();
303        let introspection = CollectionIntrospection::new(
304            id,
305            self.introspection_tx.clone(),
306            as_of.clone(),
307            storage_sink,
308            initial_as_of,
309            refresh_schedule,
310            dependency_ids,
311        );
312        let mut state = CollectionState::new(
313            id,
314            as_of.clone(),
315            shared,
316            storage_dependencies,
317            compute_dependencies,
318            Arc::clone(&self.read_hold_tx),
319            introspection,
320        );
321        state.target_replica = target_replica;
322        // If the collection is write-only, clear its read policy to reflect that.
323        if write_only {
324            state.read_policy = None;
325        }
326
327        if let Some(previous) = self.collections.insert(id, state) {
328            panic!("attempt to add a collection with existing ID {id} (previous={previous:?}");
329        }
330
331        // Add per-replica collection state.
332        for replica in self.replicas.values_mut() {
333            if target_replica.is_some_and(|id| id != replica.id) {
334                continue;
335            }
336            replica.add_collection(id, as_of.clone(), replica_input_read_holds.clone());
337        }
338    }
339
340    fn remove_collection(&mut self, id: GlobalId) {
341        // Remove per-replica collection state.
342        for replica in self.replicas.values_mut() {
343            replica.remove_collection(id);
344        }
345
346        // Remove global collection state.
347        self.collections.remove(&id);
348    }
349
350    fn add_replica_state(
351        &mut self,
352        id: ReplicaId,
353        client: ReplicaClient,
354        config: ReplicaConfig,
355        epoch: u64,
356    ) -> Result<(), read_holds::ReadHoldIssuerHungUp> {
357        let log_ids: BTreeSet<_> = config.logging.index_logs.values().copied().collect();
358
359        let metrics = self.metrics.for_replica(id);
360        let mut replica = ReplicaState::new(
361            id,
362            client,
363            config,
364            metrics,
365            self.introspection_tx.clone(),
366            epoch,
367        );
368
369        // Add per-replica collection state.
370        let mut shutdown_input = None;
371        for (collection_id, collection) in &self.collections {
372            // Skip log collections not maintained by this replica,
373            // and collections targeted at a different replica.
374            if (collection.log_collection && !log_ids.contains(collection_id))
375                || collection.target_replica.is_some_and(|rid| rid != id)
376            {
377                continue;
378            }
379
380            let as_of = if collection.log_collection {
381                // For log collections, we don't send a `CreateDataflow` command to the replica, so
382                // it doesn't know which as-of the controler chose and defaults to the minimum
383                // frontier instead. We need to initialize the controller-side tracking with the
384                // same frontier, to avoid observing regressions in the reported frontiers.
385                Antichain::from_elem(Timestamp::MIN)
386            } else {
387                collection.read_frontier().to_owned()
388            };
389
390            // Cloning a `ReadHold` fails when its issuer has hung up. For these holds the issuer
391            // is the `StorageCollections`, which doesn't hang up as long as the `Instance` exists,
392            // except during process shutdown, when the tokio runtime drops tasks in arbitrary
393            // order. In that case there is no way of correctly initializing the per-replica
394            // collection state, so we give up. We still add the replica itself, to keep the
395            // bookkeeping consistent with the controller's, and then signal the unrecoverable
396            // error to the caller, which shuts the instance down.
397            let mut input_read_holds = Vec::with_capacity(collection.storage_dependencies.len());
398            let mut hung_up = Vec::new();
399            for hold in collection.storage_dependencies.values() {
400                match hold.try_clone() {
401                    Ok(hold) => input_read_holds.push(hold),
402                    Err(read_holds::ReadHoldIssuerHungUp(input_id)) => hung_up.push(input_id),
403                }
404            }
405            if !hung_up.is_empty() {
406                tracing::error!(
407                    replica_id = %id,
408                    %collection_id,
409                    ?hung_up,
410                    "giving up on adding replica collections: storage read hold issuers hung \
411                     up, the process is shutting down",
412                );
413                shutdown_input = hung_up.into_iter().next();
414                break;
415            }
416
417            replica.add_collection(*collection_id, as_of, input_read_holds);
418        }
419
420        self.replicas.insert(id, replica);
421
422        match shutdown_input {
423            Some(input_id) => Err(read_holds::ReadHoldIssuerHungUp(input_id)),
424            None => Ok(()),
425        }
426    }
427
428    /// Enqueue the given response for delivery to the controller clients.
429    fn deliver_response(&self, response: ComputeControllerResponse) {
430        // Failure to send means the `ComputeController` has been dropped and doesn't care about
431        // responses anymore.
432        let _ = self.response_tx.send(response);
433    }
434
435    /// Enqueue the given introspection updates for recording.
436    fn deliver_introspection_updates(&self, type_: IntrospectionType, updates: Vec<(Row, Diff)>) {
437        // Failure to send means the `ComputeController` has been dropped and doesn't care about
438        // introspection updates anymore.
439        let _ = self.introspection_tx.send((type_, updates));
440    }
441
442    /// Returns whether the identified replica exists.
443    fn replica_exists(&self, id: ReplicaId) -> bool {
444        self.replicas.contains_key(&id)
445    }
446
447    /// Return the IDs of pending peeks targeting the specified replica.
448    fn peeks_targeting(&self, replica_id: ReplicaId) -> impl Iterator<Item = (Uuid, &PendingPeek)> {
449        self.peeks.iter().filter_map(move |(uuid, peek)| {
450            if peek.target_replica == Some(replica_id) {
451                Some((*uuid, peek))
452            } else {
453                None
454            }
455        })
456    }
457
458    /// Return the IDs of in-progress subscribes targeting the specified replica.
459    fn subscribes_targeting(&self, replica_id: ReplicaId) -> impl Iterator<Item = GlobalId> + '_ {
460        self.subscribes.keys().copied().filter(move |id| {
461            let collection = self.expect_collection(*id);
462            collection.target_replica == Some(replica_id)
463        })
464    }
465
466    /// Update introspection with the current collection frontiers.
467    ///
468    /// We could also do this directly in response to frontier changes, but doing it periodically
469    /// lets us avoid emitting some introspection updates that can be consolidated (e.g. a write
470    /// frontier updated immediately followed by a read frontier update).
471    ///
472    /// This method is invoked by `ComputeController::maintain`, which we expect to be called once
473    /// per second during normal operation.
474    fn update_frontier_introspection(&mut self) {
475        for collection in self.collections.values_mut() {
476            collection
477                .introspection
478                .observe_frontiers(&collection.read_frontier(), &collection.write_frontier());
479        }
480
481        for replica in self.replicas.values_mut() {
482            for collection in replica.collections.values_mut() {
483                collection
484                    .introspection
485                    .observe_frontier(&collection.write_frontier);
486            }
487        }
488    }
489
490    /// Refresh the controller state metrics for this instance.
491    ///
492    /// We could also do state metric updates directly in response to state changes, but that would
493    /// mean littering the code with metric update calls. Encapsulating state metric maintenance in
494    /// a single method is less noisy.
495    ///
496    /// This method is invoked by `ComputeController::maintain`, which we expect to be called once
497    /// per second during normal operation.
498    fn refresh_state_metrics(&self) {
499        let unscheduled_collections_count =
500            self.collections.values().filter(|c| !c.scheduled).count();
501        let connected_replica_count = self
502            .replicas
503            .values()
504            .filter(|r| r.client.is_connected())
505            .count();
506
507        self.metrics
508            .replica_count
509            .set(u64::cast_from(self.replicas.len()));
510        self.metrics
511            .collection_count
512            .set(u64::cast_from(self.collections.len()));
513        self.metrics
514            .collection_unscheduled_count
515            .set(u64::cast_from(unscheduled_collections_count));
516        self.metrics
517            .peek_count
518            .set(u64::cast_from(self.peeks.len()));
519        self.metrics
520            .subscribe_count
521            .set(u64::cast_from(self.subscribes.len()));
522        self.metrics
523            .copy_to_count
524            .set(u64::cast_from(self.copy_tos.len()));
525        self.metrics
526            .connected_replica_count
527            .set(u64::cast_from(connected_replica_count));
528    }
529
530    /// Refresh the wallclock lag introspection and metrics with the current lag values.
531    ///
532    /// This method produces wallclock lag metrics of two different shapes:
533    ///
534    /// * Histories: For each replica and each collection, we measure the lag of the write frontier
535    ///   behind the wallclock time every second. Every minute we emit the maximum lag observed
536    ///   over the last minute, together with the current time.
537    /// * Histograms: For each collection, we measure the lag of the write frontier behind
538    ///   wallclock time every second. Every minute we emit all lags observed over the last minute,
539    ///   together with the current histogram period.
540    ///
541    /// Histories are emitted to both Mz introspection and Prometheus, histograms only to
542    /// introspection. We treat lags of unreadable collections (i.e. collections that contain no
543    /// readable times) as undefined and set them to NULL in introspection and `u64::MAX` in
544    /// Prometheus.
545    ///
546    /// This method is invoked by `ComputeController::maintain`, which we expect to be called once
547    /// per second during normal operation.
548    fn refresh_wallclock_lag(&mut self) {
549        let frontier_lag = |frontier: &Antichain<Timestamp>| match frontier.as_option() {
550            Some(ts) => (self.wallclock_lag)(ts.clone()),
551            None => Duration::ZERO,
552        };
553
554        let now_ms = (self.now)();
555        let histogram_period = WallclockLagHistogramPeriod::from_epoch_millis(now_ms, &self.dyncfg);
556        let histogram_labels = match &self.workload_class {
557            Some(wc) => [("workload_class", wc.clone())].into(),
558            None => BTreeMap::new(),
559        };
560
561        // For collections that sink into storage, we need to ask the storage controller to know
562        // whether they're currently readable.
563        let readable_storage_collections: BTreeSet<_> = self
564            .collections
565            .keys()
566            .filter_map(|id| {
567                let frontiers = self.storage_collections.collection_frontiers(*id).ok()?;
568                PartialOrder::less_than(&frontiers.read_capabilities, &frontiers.write_frontier)
569                    .then_some(*id)
570            })
571            .collect();
572
573        // First, iterate over all collections and collect histogram measurements.
574        for (id, collection) in &mut self.collections {
575            let write_frontier = collection.write_frontier();
576            let readable = if self.storage_collections.check_exists(*id).is_ok() {
577                readable_storage_collections.contains(id)
578            } else {
579                PartialOrder::less_than(&collection.read_frontier(), &write_frontier)
580            };
581
582            if let Some(stash) = &mut collection.wallclock_lag_histogram_stash {
583                let bucket = if readable {
584                    let lag = frontier_lag(&write_frontier);
585                    let lag = lag.as_secs().next_power_of_two();
586                    WallclockLag::Seconds(lag)
587                } else {
588                    WallclockLag::Undefined
589                };
590
591                let key = (histogram_period, bucket, histogram_labels.clone());
592                *stash.entry(key).or_default() += Diff::ONE;
593            }
594        }
595
596        // Second, iterate over all per-replica collections and collect history measurements.
597        for replica in self.replicas.values_mut() {
598            for (id, collection) in &mut replica.collections {
599                // A per-replica collection is considered readable in the context of lag
600                // measurement if either:
601                //  (a) it sinks into a storage collection that is readable
602                //  (b) it is hydrated
603                let readable = readable_storage_collections.contains(id) || collection.hydrated();
604
605                let lag = if readable {
606                    let lag = frontier_lag(&collection.write_frontier);
607                    WallclockLag::Seconds(lag.as_secs())
608                } else {
609                    WallclockLag::Undefined
610                };
611
612                if let Some(wallclock_lag_max) = &mut collection.wallclock_lag_max {
613                    *wallclock_lag_max = (*wallclock_lag_max).max(lag);
614                }
615
616                if let Some(metrics) = &mut collection.metrics {
617                    // No way to specify values as undefined in Prometheus metrics, so we use the
618                    // maximum value instead.
619                    let secs = lag.unwrap_seconds_or(u64::MAX);
620                    metrics.wallclock_lag.observe(secs);
621                };
622            }
623        }
624
625        // Record lags to persist, if it's time.
626        self.maybe_record_wallclock_lag();
627    }
628
629    /// Produce new wallclock lag introspection updates, provided enough time has passed since the
630    /// last recording.
631    //
632    /// We emit new introspection updates if the system time has passed into a new multiple of the
633    /// recording interval (typically 1 minute) since the last refresh. The storage controller uses
634    /// the same approach, ensuring that both controllers commit their lags at roughly the same
635    /// time, avoiding confusion caused by inconsistencies.
636    fn maybe_record_wallclock_lag(&mut self) {
637        if self.read_only {
638            return;
639        }
640
641        let duration_trunc = |datetime: DateTime<_>, interval| {
642            let td = TimeDelta::from_std(interval).ok()?;
643            datetime.duration_trunc(td).ok()
644        };
645
646        let interval = WALLCLOCK_LAG_RECORDING_INTERVAL.get(&self.dyncfg);
647        let now_dt = mz_ore::now::to_datetime((self.now)());
648        let now_trunc = duration_trunc(now_dt, interval).unwrap_or_else(|| {
649            soft_panic_or_log!("excessive wallclock lag recording interval: {interval:?}");
650            let default = WALLCLOCK_LAG_RECORDING_INTERVAL.default();
651            duration_trunc(now_dt, *default).unwrap()
652        });
653        if now_trunc <= self.wallclock_lag_last_recorded {
654            return;
655        }
656
657        let now_ts: CheckedTimestamp<_> = now_trunc.try_into().expect("must fit");
658
659        let mut history_updates = Vec::new();
660        for (replica_id, replica) in &mut self.replicas {
661            for (collection_id, collection) in &mut replica.collections {
662                let Some(wallclock_lag_max) = &mut collection.wallclock_lag_max else {
663                    continue;
664                };
665
666                let max_lag = std::mem::replace(wallclock_lag_max, WallclockLag::MIN);
667                let row = Row::pack_slice(&[
668                    Datum::String(&collection_id.to_string()),
669                    Datum::String(&replica_id.to_string()),
670                    max_lag.into_interval_datum(),
671                    Datum::TimestampTz(now_ts),
672                ]);
673                history_updates.push((row, Diff::ONE));
674            }
675        }
676        if !history_updates.is_empty() {
677            self.deliver_introspection_updates(
678                IntrospectionType::WallclockLagHistory,
679                history_updates,
680            );
681        }
682
683        let mut histogram_updates = Vec::new();
684        let mut row_buf = Row::default();
685        for (collection_id, collection) in &mut self.collections {
686            let Some(stash) = &mut collection.wallclock_lag_histogram_stash else {
687                continue;
688            };
689
690            for ((period, lag, labels), count) in std::mem::take(stash) {
691                let mut packer = row_buf.packer();
692                packer.extend([
693                    Datum::TimestampTz(period.start),
694                    Datum::TimestampTz(period.end),
695                    Datum::String(&collection_id.to_string()),
696                    lag.into_uint64_datum(),
697                ]);
698                let labels = labels.iter().map(|(k, v)| (*k, Datum::String(v)));
699                packer.push_dict(labels);
700
701                histogram_updates.push((row_buf.clone(), count));
702            }
703        }
704        if !histogram_updates.is_empty() {
705            self.deliver_introspection_updates(
706                IntrospectionType::WallclockLagHistogram,
707                histogram_updates,
708            );
709        }
710
711        self.wallclock_lag_last_recorded = now_trunc;
712    }
713
714    /// Returns `true` if the given collection is hydrated on at least one
715    /// replica.
716    ///
717    /// This also returns `true` in case this cluster does not have any
718    /// replicas that host the given collection.
719    #[mz_ore::instrument(level = "debug")]
720    pub fn collection_hydrated(&self, collection_id: GlobalId) -> Result<bool, CollectionMissing> {
721        let mut hosting_replicas = self.replicas_hosting(collection_id)?.peekable();
722        if hosting_replicas.peek().is_none() {
723            return Ok(true);
724        }
725        for replica_state in hosting_replicas {
726            if replica_state.expect_collection(collection_id).hydrated() {
727                return Ok(true);
728            }
729        }
730
731        Ok(false)
732    }
733
734    /// Returns `true` if each non-transient, non-excluded collection is *ready* on at
735    /// least one of the target replicas.
736    ///
737    /// A target must be hydrated and, when `allowed_lag` is `Some`, its output
738    /// frontier must be within that allowance of the furthest output frontier
739    /// among the collection's hosting `reference_replica_ids`. No reference
740    /// replicas means hydration suffices. An empty reference frontier requires
741    /// the target to complete too. `None` checks only hydration.
742    ///
743    /// Hydration alone only guarantees output past the installation as-of, not
744    /// that the replica has replayed subsequent updates. Output frontiers measure
745    /// that replay. MV write frontiers can instead reflect another replica's
746    /// writes to the shared shard, or jump ahead to the next `REFRESH` time.
747    ///
748    /// Zero-replica clusters return `true`.
749    #[mz_ore::instrument(level = "debug")]
750    pub fn collections_ready_on_replicas(
751        &self,
752        target_replica_ids: Option<Vec<ReplicaId>>,
753        exclude_collections: &BTreeSet<GlobalId>,
754        allowed_lag: Option<Timestamp>,
755        reference_replica_ids: &BTreeSet<ReplicaId>,
756    ) -> Result<bool, HydrationCheckBadTarget> {
757        if self.replicas.is_empty() {
758            return Ok(true);
759        }
760        let target_replicas: BTreeSet<ReplicaId> = self
761            .replicas
762            .keys()
763            .filter_map(|id| match target_replica_ids {
764                None => Some(id.clone()),
765                Some(ref ids) if ids.contains(id) => Some(id.clone()),
766                Some(_) => None,
767            })
768            .collect();
769        if let Some(targets) = target_replica_ids {
770            if target_replicas.is_empty() {
771                return Err(HydrationCheckBadTarget(targets));
772            }
773        }
774
775        let mut unhydrated = BTreeSet::new();
776        let mut lagging_ticks = BTreeMap::new();
777        let mut awaiting_completion = BTreeSet::new();
778        for (id, _collection) in self.collections_iter() {
779            if id.is_transient() || exclude_collections.contains(&id) {
780                continue;
781            }
782
783            // `replicas_hosting` cannot fail here because `collections_iter`
784            // only yields collections that exist.
785            let replicas = self
786                .replicas_hosting(id)
787                .expect("collection must exist")
788                .map(|replica| (replica.id, replica.expect_collection(id)));
789
790            match classify_collection_readiness(
791                replicas,
792                &target_replicas,
793                reference_replica_ids,
794                allowed_lag,
795            ) {
796                CollectionReadiness::Ready => {}
797                // We collect all not-ready collections instead of breaking out
798                // early, so that the log below names every collection the caller
799                // is waiting on, and why.
800                CollectionReadiness::Lagging { lag: Some(lag) } => {
801                    lagging_ticks.insert(id, lag);
802                }
803                CollectionReadiness::Lagging { lag: None } => {
804                    awaiting_completion.insert(id);
805                }
806                CollectionReadiness::Unhydrated => {
807                    unhydrated.insert(id);
808                }
809            }
810        }
811
812        let ready =
813            unhydrated.is_empty() && lagging_ticks.is_empty() && awaiting_completion.is_empty();
814        if !ready {
815            // Callers poll this on the cluster controller's reconcile tick,
816            // which tests turn down to milliseconds, so this is deliberately
817            // one line per call rather than one per collection.
818            tracing::info!(
819                replicas = ?target_replicas,
820                reference = ?reference_replica_ids,
821                unhydrated = ?unhydrated,
822                ?lagging_ticks,
823                ?awaiting_completion,
824                ?allowed_lag,
825                "collections are not ready on any target replica",
826            );
827        }
828
829        Ok(ready)
830    }
831
832    /// Clean up collection state that is not needed anymore.
833    ///
834    /// Three conditions need to be true before we can remove state for a collection:
835    ///
836    ///  1. A client must have explicitly dropped the collection. If that is not the case, clients
837    ///     can still reasonably assume that the controller knows about the collection and can
838    ///     answer queries about it.
839    ///  2. There must be no outstanding read capabilities on the collection. As long as someone
840    ///     still holds read capabilities on a collection, we need to keep it around to be able
841    ///     to properly handle downgrading of said capabilities.
842    ///  3. All replica frontiers for the collection must have advanced to the empty frontier.
843    ///     Advancement to the empty frontiers signals that replicas are done computing the
844    ///     collection and that they won't send more `ComputeResponse`s for it. As long as we might
845    ///     receive responses for a collection we want to keep it around to be able to validate and
846    ///     handle these responses.
847    fn cleanup_collections(&mut self) {
848        let to_remove: Vec<_> = self
849            .collections_iter()
850            .filter(|(id, collection)| {
851                collection.dropped
852                    && collection.shared.lock_read_capabilities(|c| c.is_empty())
853                    && self
854                        .replicas
855                        .values()
856                        .all(|r| r.collection_frontiers_empty(*id))
857            })
858            .map(|(id, _collection)| id)
859            .collect();
860
861        for id in to_remove {
862            self.remove_collection(id);
863        }
864    }
865
866    /// Returns the state of the [`Instance`] formatted as JSON.
867    ///
868    /// The returned value is not guaranteed to be stable and may change at any point in time.
869    #[mz_ore::instrument(level = "debug")]
870    pub fn dump(&self) -> Result<serde_json::Value, anyhow::Error> {
871        // Note: We purposefully use the `Debug` formatting for the value of all fields in the
872        // returned object as a tradeoff between usability and stability. `serde_json` will fail
873        // to serialize an object if the keys aren't strings, so `Debug` formatting the values
874        // prevents a future unrelated change from silently breaking this method.
875
876        // Destructure `self` here so we don't forget to consider dumping newly added fields.
877        let Self {
878            build_info: _,
879            storage_collections: _,
880            peek_stash_persist_location: _,
881            initialized,
882            read_only,
883            workload_class,
884            replicas,
885            replica_dyncfg_overrides: _,
886            collections,
887            log_sources: _,
888            peeks,
889            subscribes,
890            copy_tos,
891            history: _,
892            command_rx: _,
893            response_tx: _,
894            introspection_tx: _,
895            metrics: _,
896            dyncfg: _,
897            now: _,
898            wallclock_lag: _,
899            wallclock_lag_last_recorded,
900            read_hold_tx: _,
901            replica_tx: _,
902            replica_rx: _,
903        } = self;
904
905        let replicas: BTreeMap<_, _> = replicas
906            .iter()
907            .map(|(id, replica)| Ok((id.to_string(), replica.dump()?)))
908            .collect::<Result<_, anyhow::Error>>()?;
909        let collections: BTreeMap<_, _> = collections
910            .iter()
911            .map(|(id, collection)| (id.to_string(), format!("{collection:?}")))
912            .collect();
913        let peeks: BTreeMap<_, _> = peeks
914            .iter()
915            .map(|(uuid, peek)| (uuid.to_string(), format!("{peek:?}")))
916            .collect();
917        let subscribes: BTreeMap<_, _> = subscribes
918            .iter()
919            .map(|(id, subscribe)| (id.to_string(), format!("{subscribe:?}")))
920            .collect();
921        let copy_tos: Vec<_> = copy_tos.iter().map(|id| id.to_string()).collect();
922        let wallclock_lag_last_recorded = format!("{wallclock_lag_last_recorded:?}");
923
924        Ok(serde_json::json!({
925            "initialized": initialized,
926            "read_only": read_only,
927            "workload_class": workload_class,
928            "replicas": replicas,
929            "collections": collections,
930            "peeks": peeks,
931            "subscribes": subscribes,
932            "copy_tos": copy_tos,
933            "wallclock_lag_last_recorded": wallclock_lag_last_recorded,
934        }))
935    }
936
937    /// Reports the current write frontier for the identified compute collection.
938    pub(super) fn collection_write_frontier(
939        &self,
940        id: GlobalId,
941    ) -> Result<Antichain<Timestamp>, CollectionMissing> {
942        Ok(self.collection(id)?.write_frontier())
943    }
944}
945
946impl Instance {
947    pub(super) fn new(
948        build_info: &'static BuildInfo,
949        storage: StorageCollections,
950        peek_stash_persist_location: PersistLocation,
951        arranged_logs: Vec<(LogVariant, GlobalId, SharedCollectionState)>,
952        metrics: InstanceMetrics,
953        now: NowFn,
954        wallclock_lag: WallclockLagFn<Timestamp>,
955        dyncfg: Arc<ConfigSet>,
956        command_rx: mpsc::UnboundedReceiver<Command>,
957        response_tx: mpsc::UnboundedSender<ComputeControllerResponse>,
958        read_hold_tx: read_holds::ChangeTx,
959        introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
960        read_only: bool,
961    ) -> Self {
962        let mut collections = BTreeMap::new();
963        let mut log_sources = BTreeMap::new();
964        for (log, id, shared) in arranged_logs {
965            let collection = CollectionState::new_log_collection(
966                id,
967                shared,
968                Arc::clone(&read_hold_tx),
969                introspection_tx.clone(),
970            );
971            collections.insert(id, collection);
972            log_sources.insert(log, id);
973        }
974
975        let history = ComputeCommandHistory::new(metrics.for_history());
976
977        let send_count = metrics.response_send_count.clone();
978        let recv_count = metrics.response_recv_count.clone();
979        let (replica_tx, replica_rx) = instrumented_unbounded_channel(send_count, recv_count);
980
981        let now_dt = mz_ore::now::to_datetime(now());
982
983        Self {
984            build_info,
985            storage_collections: storage,
986            peek_stash_persist_location,
987            initialized: false,
988            read_only,
989            workload_class: None,
990            replicas: Default::default(),
991            replica_dyncfg_overrides: Default::default(),
992            collections,
993            log_sources,
994            peeks: Default::default(),
995            subscribes: Default::default(),
996            copy_tos: Default::default(),
997            history,
998            command_rx,
999            response_tx,
1000            introspection_tx,
1001            metrics,
1002            dyncfg,
1003            now,
1004            wallclock_lag,
1005            wallclock_lag_last_recorded: now_dt,
1006            read_hold_tx,
1007            replica_tx,
1008            replica_rx,
1009        }
1010    }
1011
1012    pub(super) async fn run(mut self) {
1013        self.send(ComputeCommand::Hello {
1014            // The nonce is protocol iteration-specific and will be set in
1015            // `ReplicaTask::specialize_command`.
1016            nonce: Uuid::default(),
1017        });
1018
1019        let instance_config = InstanceConfig {
1020            peek_stash_persist_location: self.peek_stash_persist_location.clone(),
1021            // The remaining fields are replica-specific and will be set in
1022            // `ReplicaTask::specialize_command` (logging, expiration, dictionary compression) and
1023            // `Instance::specialize_command_for_replica` (the initial config snapshot).
1024            logging: Default::default(),
1025            expiration_offset: Default::default(),
1026            arrangement_dictionary_compression: Default::default(),
1027            initial_config: Default::default(),
1028        };
1029
1030        self.send(ComputeCommand::CreateInstance(Box::new(instance_config)));
1031
1032        loop {
1033            tokio::select! {
1034                command = self.command_rx.recv() => match command {
1035                    Some(cmd) => cmd(&mut self),
1036                    None => break,
1037                },
1038                response = self.replica_rx.recv() => match response {
1039                    Some(response) => self.handle_response(response),
1040                    None => unreachable!("self owns a sender side of the channel"),
1041                }
1042            }
1043        }
1044    }
1045
1046    /// Update instance configuration.
1047    #[mz_ore::instrument(level = "debug")]
1048    pub fn update_configuration(&mut self, config_params: ComputeParameters) {
1049        if let Some(workload_class) = &config_params.workload_class {
1050            self.workload_class = workload_class.clone();
1051        }
1052
1053        let command = ComputeCommand::UpdateConfiguration(Box::new(config_params));
1054        self.send(command);
1055    }
1056
1057    /// Marks the end of any initialization commands.
1058    ///
1059    /// Intended to be called by `Controller`, rather than by other code.
1060    /// Calling this method repeatedly has no effect.
1061    #[mz_ore::instrument(level = "debug")]
1062    pub fn initialization_complete(&mut self) {
1063        // The compute protocol requires that `InitializationComplete` is sent only once.
1064        if !self.initialized {
1065            self.send(ComputeCommand::InitializationComplete);
1066            self.initialized = true;
1067        }
1068    }
1069
1070    /// Allows collections to affect writes to external systems (persist).
1071    ///
1072    /// Calling this method repeatedly has no effect.
1073    #[mz_ore::instrument(level = "debug")]
1074    pub fn allow_writes(&mut self, collection_id: GlobalId) -> Result<(), CollectionMissing> {
1075        let collection = self.collection_mut(collection_id)?;
1076
1077        // Do not send redundant allow-writes commands.
1078        if !collection.read_only {
1079            return Ok(());
1080        }
1081
1082        // Don't send allow-writes for collections that are not installed.
1083        let as_of = collection.read_frontier();
1084
1085        // If the collection has an empty `as_of`, it was either never installed on the replica or
1086        // has since been dropped. In either case the replica does not expect any commands for it.
1087        if as_of.is_empty() {
1088            return Ok(());
1089        }
1090
1091        collection.read_only = false;
1092        self.send(ComputeCommand::AllowWrites(collection_id));
1093
1094        Ok(())
1095    }
1096
1097    /// Shut down this instance.
1098    ///
1099    /// This method asserts that the instance has no replicas left. It exists to help
1100    /// us find bugs where the client drops a compute instance that still has replicas
1101    /// installed, and later assumes that said replicas still exist.
1102    ///
1103    /// # Panics
1104    ///
1105    /// Soft-panics if the compute instance still has active replicas.
1106    #[mz_ore::instrument(level = "debug")]
1107    pub fn shutdown(&mut self) {
1108        // Taking the `command_rx` ensures that the [`Instance::run`] loop terminates.
1109        let (_tx, rx) = mpsc::unbounded_channel();
1110        self.command_rx = rx;
1111
1112        let stray_replicas: Vec<_> = self.replicas.keys().collect();
1113        soft_assert_or_log!(
1114            stray_replicas.is_empty(),
1115            "dropped instance still has provisioned replicas: {stray_replicas:?}",
1116        );
1117    }
1118
1119    /// Terminate the [`Instance::run`] loop, causing the instance task to shut down.
1120    ///
1121    /// Unlike [`Instance::shutdown`], this does not assert that the instance has no replicas
1122    /// left. We use it to react to unrecoverable errors that can only occur during process
1123    /// shutdown, such as a storage read hold issuer hanging up while we rehydrate a replica.
1124    fn initiate_shutdown(&mut self) {
1125        // Replacing `command_rx` with a fresh, sender-less channel makes the next `recv` in
1126        // [`Instance::run`] return `None`, terminating the loop.
1127        let (_tx, rx) = mpsc::unbounded_channel();
1128        self.command_rx = rx;
1129    }
1130
1131    /// Sends a command to replicas of this instance.
1132    #[mz_ore::instrument(level = "debug")]
1133    fn send(&mut self, cmd: ComputeCommand) {
1134        // Record the command so that new replicas can be brought up to speed.
1135        // We record the *base* (un-specialized) command, so that the per-replica
1136        // dyncfg overrides are re-applied at replay time in `add_replica` rather
1137        // than baked into the shared history.
1138        self.history.push(cmd.clone());
1139
1140        let target_replica = self.target_replica(&cmd);
1141
1142        // Borrow the overrides and dyncfg separately from `self.replicas` so the per-replica
1143        // specialization below does not conflict with the mutable replica borrow.
1144        let overrides = &self.replica_dyncfg_overrides;
1145        let dyncfg = &self.dyncfg;
1146
1147        if let Some(rid) = target_replica {
1148            if let Some(replica) = self.replicas.get_mut(&rid) {
1149                let cmd = Self::specialize_command_for_replica(cmd, rid, overrides, dyncfg);
1150                let _ = replica.client.send(cmd);
1151            }
1152        } else {
1153            for (rid, replica) in self.replicas.iter_mut() {
1154                let cmd =
1155                    Self::specialize_command_for_replica(cmd.clone(), *rid, overrides, dyncfg);
1156                let _ = replica.client.send(cmd);
1157            }
1158        }
1159    }
1160
1161    /// Specializes a command for a specific replica by merging that replica's dyncfg override into
1162    /// its configuration. For `UpdateConfiguration` the override is merged into the update. For
1163    /// `CreateInstance` the current dyncfg, with the override applied on top, is captured as the
1164    /// initial config snapshot. All other commands are returned unchanged.
1165    ///
1166    /// The snapshot is built here, rather than baked into the history, so it reflects the dyncfg
1167    /// and override values current at the time the command is sent or replayed to the replica.
1168    fn specialize_command_for_replica(
1169        mut cmd: ComputeCommand,
1170        replica_id: ReplicaId,
1171        overrides: &BTreeMap<ReplicaId, ConfigUpdates>,
1172        dyncfg: &ConfigSet,
1173    ) -> ComputeCommand {
1174        let over = overrides.get(&replica_id);
1175        match &mut cmd {
1176            ComputeCommand::UpdateConfiguration(params) => {
1177                if let Some(over) = over
1178                    && !over.updates.is_empty()
1179                {
1180                    params.dyncfg_updates.extend(over.clone());
1181                }
1182            }
1183            ComputeCommand::CreateInstance(config) => {
1184                let mut initial = ConfigUpdates::from(dyncfg);
1185                if let Some(over) = over {
1186                    initial.extend(over.clone());
1187                }
1188                config.initial_config = initial;
1189            }
1190            _ => {}
1191        }
1192        cmd
1193    }
1194
1195    /// Replaces the per-replica dyncfg overrides. Callers should follow this
1196    /// with a configuration push (e.g. `update_configuration`) so that existing
1197    /// replicas observe the new overrides.
1198    pub(super) fn update_replica_dyncfg_overrides(
1199        &mut self,
1200        overrides: BTreeMap<ReplicaId, ConfigUpdates>,
1201    ) {
1202        self.replica_dyncfg_overrides = overrides;
1203    }
1204
1205    /// Determine the target replica for a compute command. Retrieves the
1206    /// collection named by the command, and returns the target replica if
1207    /// it is set, and None if not set, or the command doesn't name a collection.
1208    ///
1209    /// Panics if a create-dataflow command names collections that have different
1210    /// target replicas. It is an error to construct such an object and would
1211    /// indicate a bug in [`Self::create_dataflow`].
1212    fn target_replica(&self, cmd: &ComputeCommand) -> Option<ReplicaId> {
1213        match &cmd {
1214            ComputeCommand::Schedule(id)
1215            | ComputeCommand::AllowWrites(id)
1216            | ComputeCommand::AllowCompaction { id, .. } => {
1217                self.expect_collection(*id).target_replica
1218            }
1219            ComputeCommand::CreateDataflow(desc) => {
1220                let mut target_replica = None;
1221                for id in desc.export_ids() {
1222                    if let Some(replica) = self.expect_collection(id).target_replica {
1223                        if target_replica.is_some() {
1224                            assert_eq!(target_replica, Some(replica));
1225                        }
1226                        target_replica = Some(replica);
1227                    }
1228                }
1229                target_replica
1230            }
1231            // Skip Peek as we don't allow replica-targeted indexes.
1232            ComputeCommand::Peek(_)
1233            | ComputeCommand::Hello { .. }
1234            | ComputeCommand::CreateInstance(_)
1235            | ComputeCommand::InitializationComplete
1236            | ComputeCommand::UpdateConfiguration(_)
1237            | ComputeCommand::CancelPeek { .. } => None,
1238        }
1239    }
1240
1241    /// Add a new instance replica, by ID.
1242    #[mz_ore::instrument(level = "debug")]
1243    pub fn add_replica(
1244        &mut self,
1245        id: ReplicaId,
1246        mut config: ReplicaConfig,
1247        epoch: Option<u64>,
1248    ) -> Result<(), ReplicaExists> {
1249        if self.replica_exists(id) {
1250            return Err(ReplicaExists(id));
1251        }
1252
1253        config.logging.index_logs = self.log_sources.clone();
1254
1255        let epoch = epoch.unwrap_or(1);
1256        let metrics = self.metrics.for_replica(id);
1257        let client = ReplicaClient::spawn(
1258            id,
1259            self.build_info,
1260            config.clone(),
1261            epoch,
1262            metrics.clone(),
1263            Arc::clone(&self.dyncfg),
1264            self.replica_tx.clone(),
1265        );
1266
1267        // Take this opportunity to clean up the history we should present.
1268        self.history.reduce();
1269
1270        // Advance the uppers of source imports
1271        self.history.update_source_uppers(&self.storage_collections);
1272
1273        // Replay the commands at the client, creating new dataflow identifiers.
1274        for command in self.history.iter() {
1275            // Skip `CreateDataflow` commands targeted at different replicas.
1276            if let Some(target_replica) = self.target_replica(command)
1277                && target_replica != id
1278            {
1279                continue;
1280            }
1281
1282            // Re-apply this replica's dyncfg override to replayed config commands, and rebuild the
1283            // create-instance snapshot from the current dyncfg.
1284            let command = Self::specialize_command_for_replica(
1285                command.clone(),
1286                id,
1287                &self.replica_dyncfg_overrides,
1288                &self.dyncfg,
1289            );
1290            if client.send(command).is_err() {
1291                // We swallow the error here. On the next send, we will fail again, and
1292                // restart the connection as well as this rehydration.
1293                tracing::warn!("Replica {:?} connection terminated during hydration", id);
1294                break;
1295            }
1296        }
1297
1298        // Add replica to tracked state.
1299        if self.add_replica_state(id, client, config, epoch).is_err() {
1300            // A storage read hold issuer hung up, which only happens during process shutdown.
1301            // There is no way to correctly bring up the replica anymore, so we shut the instance
1302            // down instead of running on with half-initialized replica state. `add_replica_state`
1303            // has already logged the details and inserted the replica to keep our bookkeeping
1304            // consistent with the controller's.
1305            self.initiate_shutdown();
1306        }
1307
1308        Ok(())
1309    }
1310
1311    /// Remove an existing instance replica, by ID.
1312    #[mz_ore::instrument(level = "debug")]
1313    pub fn remove_replica(&mut self, id: ReplicaId) -> Result<(), ReplicaMissing> {
1314        let replica = self.replicas.remove(&id).ok_or(ReplicaMissing(id))?;
1315
1316        // The coordinator only re-pushes the override map when the scoped configuration itself
1317        // changes, so a dropped replica's entry would otherwise be retained until the next such
1318        // change.
1319        self.replica_dyncfg_overrides.remove(&id);
1320
1321        // Before dropping the replica state (and the contained input read holds), log read holds
1322        // that are the last line of defense against compaction of a dataflow's storage inputs. If
1323        // the corresponding global read hold has already been released, dropping the per-replica
1324        // read hold will allow compaction, which can cause the replica to panic trying to install
1325        // the dataflow.
1326        //
1327        // This exists primarily to help diagnose incidents-and-escalations#39.
1328        for (collection_id, replica_collection) in &replica.collections {
1329            let collection = self.collections.get(collection_id);
1330            for replica_hold in &replica_collection.input_read_holds {
1331                let input_id = replica_hold.id();
1332                let global_hold = collection.and_then(|c| c.storage_dependencies.get(&input_id));
1333                let unprotected = global_hold
1334                    .is_none_or(|h| PartialOrder::less_than(replica_hold.since(), h.since()));
1335                if unprotected {
1336                    tracing::warn!(
1337                        replica_id = %id,
1338                        %collection_id,
1339                        %input_id,
1340                        replica_hold_since = ?replica_hold.since(),
1341                        global_hold_since = ?global_hold.map(|h| h.since()),
1342                        "dropping per-replica read hold without equivalent global read hold",
1343                    );
1344                }
1345            }
1346        }
1347        drop(replica);
1348
1349        // Subscribes targeting this replica either won't be served anymore (if the replica is
1350        // dropped) or might produce inconsistent output (if the target collection is an
1351        // introspection index). We produce an error to inform upstream.
1352        let to_drop: Vec<_> = self.subscribes_targeting(id).collect();
1353        for subscribe_id in to_drop {
1354            let subscribe = self.subscribes.remove(&subscribe_id).unwrap();
1355            let response = ComputeControllerResponse::SubscribeResponse(
1356                subscribe_id,
1357                SubscribeBatch {
1358                    lower: subscribe.frontier.clone(),
1359                    upper: subscribe.frontier,
1360                    updates: Err(ERROR_TARGET_REPLICA_FAILED.into()),
1361                },
1362            );
1363            self.deliver_response(response);
1364        }
1365
1366        // Peeks targeting this replica might not be served anymore (if the replica is dropped).
1367        // If the replica has failed it might come back and respond to the peek later, but it still
1368        // seems like a good idea to cancel the peek to inform the caller about the failure. This
1369        // is consistent with how we handle targeted subscribes above.
1370        let mut peek_responses = Vec::new();
1371        let mut to_drop = Vec::new();
1372        for (uuid, peek) in self.peeks_targeting(id) {
1373            peek_responses.push(ComputeControllerResponse::PeekNotification(
1374                uuid,
1375                PeekNotification::Error(ERROR_TARGET_REPLICA_FAILED.into()),
1376                peek.otel_ctx.clone(),
1377            ));
1378            to_drop.push(uuid);
1379        }
1380        for response in peek_responses {
1381            self.deliver_response(response);
1382        }
1383        for uuid in to_drop {
1384            let response =
1385                PeekResponse::Error(ProtocolPeekError::unstructured(ERROR_TARGET_REPLICA_FAILED));
1386            self.finish_peek(uuid, response);
1387        }
1388
1389        // We might have a chance to forward implied capabilities and reduce the cost of bringing
1390        // up the next replica, if the dropped replica was the only one in the cluster.
1391        self.forward_implied_capabilities();
1392
1393        Ok(())
1394    }
1395
1396    /// Rehydrate the given instance replica.
1397    ///
1398    /// # Panics
1399    ///
1400    /// Panics if the specified replica does not exist.
1401    fn rehydrate_replica(&mut self, id: ReplicaId) {
1402        let config = self.replicas[&id].config.clone();
1403        let epoch = self.replicas[&id].epoch + 1;
1404
1405        self.remove_replica(id).expect("replica must exist");
1406        let result = self.add_replica(id, config, Some(epoch));
1407
1408        match result {
1409            Ok(()) => (),
1410            Err(ReplicaExists(_)) => unreachable!("replica was removed"),
1411        }
1412    }
1413
1414    /// Rehydrate any failed replicas of this instance.
1415    fn rehydrate_failed_replicas(&mut self) {
1416        let replicas = self.replicas.iter();
1417        let failed_replicas: Vec<_> = replicas
1418            .filter_map(|(id, replica)| replica.client.is_failed().then_some(*id))
1419            .collect();
1420
1421        for replica_id in failed_replicas {
1422            self.rehydrate_replica(replica_id);
1423        }
1424    }
1425
1426    /// Creates the described dataflow and initializes state for its output.
1427    ///
1428    /// This method expects a `DataflowDescription` with an `as_of` frontier specified, as well as
1429    /// for each imported collection a read hold in `import_read_holds` at at least the `as_of`.
1430    #[mz_ore::instrument(level = "debug")]
1431    pub fn create_dataflow(
1432        &mut self,
1433        dataflow: DataflowDescription<mz_compute_types::plan::LirRelationExpr, ()>,
1434        import_read_holds: Vec<ReadHold>,
1435        mut shared_collection_state: BTreeMap<GlobalId, SharedCollectionState>,
1436        target_replica: Option<ReplicaId>,
1437    ) -> Result<(), DataflowCreationError> {
1438        use DataflowCreationError::*;
1439
1440        // Validate that the target replica, if specified, exists.
1441        // A targeted dataflow is only installed on a single replica; if that
1442        // replica doesn't exist, we can't create the dataflow.
1443        if let Some(replica_id) = target_replica {
1444            if !self.replica_exists(replica_id) {
1445                return Err(ReplicaMissing(replica_id));
1446            }
1447        }
1448
1449        // Simple sanity checks around `as_of`
1450        let as_of = dataflow.as_of.as_ref().ok_or(MissingAsOf)?;
1451        if as_of.is_empty() && dataflow.subscribe_ids().next().is_some() {
1452            return Err(EmptyAsOfForSubscribe);
1453        }
1454        if as_of.is_empty() && dataflow.copy_to_ids().next().is_some() {
1455            return Err(EmptyAsOfForCopyTo);
1456        }
1457
1458        // Collect all dependencies of the dataflow, and read holds on them at the `as_of`.
1459        let mut storage_dependencies = BTreeMap::new();
1460        let mut compute_dependencies = BTreeMap::new();
1461
1462        // When we install per-replica input read holds, we cannot use the `as_of` because of
1463        // reconciliation: Existing slow replicas might be reading from the inputs at times before
1464        // the `as_of` and we would rather not crash them by allowing their inputs to compact too
1465        // far. So instead we take read holds at the least time available.
1466        let mut replica_input_read_holds = Vec::new();
1467
1468        let mut import_read_holds: BTreeMap<_, _> =
1469            import_read_holds.into_iter().map(|r| (r.id(), r)).collect();
1470
1471        for &id in dataflow.source_imports.keys() {
1472            let mut read_hold = import_read_holds.remove(&id).ok_or(ReadHoldMissing(id))?;
1473            replica_input_read_holds.push(read_hold.clone());
1474
1475            read_hold
1476                .try_downgrade(as_of.clone())
1477                .map_err(|_| ReadHoldInsufficient(id))?;
1478            storage_dependencies.insert(id, read_hold);
1479        }
1480
1481        for &id in dataflow.index_imports.keys() {
1482            let mut read_hold = import_read_holds.remove(&id).ok_or(ReadHoldMissing(id))?;
1483            read_hold
1484                .try_downgrade(as_of.clone())
1485                .map_err(|_| ReadHoldInsufficient(id))?;
1486            compute_dependencies.insert(id, read_hold);
1487        }
1488
1489        // If the `as_of` is empty, we are not going to create a dataflow, so replicas won't read
1490        // from the inputs.
1491        if as_of.is_empty() {
1492            replica_input_read_holds = Default::default();
1493        }
1494
1495        // Install collection state for each of the exports.
1496        for export_id in dataflow.export_ids() {
1497            let shared = shared_collection_state
1498                .remove(&export_id)
1499                .unwrap_or_else(|| SharedCollectionState::new(as_of.clone()));
1500            let write_only = dataflow.sink_exports.contains_key(&export_id);
1501            let storage_sink = dataflow.persist_sink_ids().any(|id| id == export_id);
1502
1503            self.add_collection(
1504                export_id,
1505                as_of.clone(),
1506                shared,
1507                storage_dependencies.clone(),
1508                compute_dependencies.clone(),
1509                replica_input_read_holds.clone(),
1510                write_only,
1511                storage_sink,
1512                dataflow.initial_storage_as_of.clone(),
1513                dataflow.refresh_schedule.clone(),
1514                target_replica,
1515            );
1516
1517            // If the export is a storage sink, we can advance its write frontier to the write
1518            // frontier of the target storage collection.
1519            if let Ok(frontiers) = self.storage_collections.collection_frontiers(export_id) {
1520                self.maybe_update_global_write_frontier(export_id, frontiers.write_frontier);
1521            }
1522        }
1523
1524        // Initialize tracking of subscribes.
1525        for subscribe_id in dataflow.subscribe_ids() {
1526            self.subscribes
1527                .insert(subscribe_id, ActiveSubscribe::default());
1528        }
1529
1530        // Initialize tracking of copy tos.
1531        for copy_to_id in dataflow.copy_to_ids() {
1532            self.copy_tos.insert(copy_to_id);
1533        }
1534
1535        // Here we augment all imported sources and all exported sinks with the appropriate
1536        // storage metadata needed by the compute instance.
1537        let mut source_imports = BTreeMap::new();
1538        for (id, import) in dataflow.source_imports {
1539            let frontiers = self
1540                .storage_collections
1541                .collection_frontiers(id)
1542                .expect("collection exists");
1543
1544            let collection_metadata = self
1545                .storage_collections
1546                .collection_metadata(id)
1547                .expect("we have a read hold on this collection");
1548
1549            let desc = SourceInstanceDesc {
1550                storage_metadata: collection_metadata.clone(),
1551                arguments: import.desc.arguments,
1552                typ: import.desc.typ.clone(),
1553            };
1554            source_imports.insert(
1555                id,
1556                mz_compute_types::dataflows::SourceImport {
1557                    desc,
1558                    monotonic: import.monotonic,
1559                    with_snapshot: import.with_snapshot,
1560                    upper: frontiers.write_frontier,
1561                },
1562            );
1563        }
1564
1565        let mut sink_exports = BTreeMap::new();
1566        for (id, se) in dataflow.sink_exports {
1567            let connection = match se.connection {
1568                ComputeSinkConnection::MaterializedView(conn) => {
1569                    let metadata = self
1570                        .storage_collections
1571                        .collection_metadata(id)
1572                        .map_err(|_| CollectionMissing(id))?
1573                        .clone();
1574                    let conn = MaterializedViewSinkConnection {
1575                        value_desc: conn.value_desc,
1576                        storage_metadata: metadata,
1577                    };
1578                    ComputeSinkConnection::MaterializedView(conn)
1579                }
1580                ComputeSinkConnection::Subscribe(conn) => ComputeSinkConnection::Subscribe(conn),
1581                ComputeSinkConnection::CopyToS3Oneshot(conn) => {
1582                    ComputeSinkConnection::CopyToS3Oneshot(conn)
1583                }
1584                ComputeSinkConnection::MetricSink(conn) => ComputeSinkConnection::MetricSink(conn),
1585            };
1586            let desc = ComputeSinkDesc {
1587                from: se.from,
1588                from_desc: se.from_desc,
1589                connection,
1590                with_snapshot: se.with_snapshot,
1591                up_to: se.up_to,
1592                non_null_assertions: se.non_null_assertions,
1593                refresh_schedule: se.refresh_schedule,
1594            };
1595            sink_exports.insert(id, desc);
1596        }
1597
1598        // Flatten the dataflow plans into the representation expected by replicas.
1599        let objects_to_build = dataflow
1600            .objects_to_build
1601            .into_iter()
1602            .map(|object| BuildDesc {
1603                id: object.id,
1604                plan: RenderPlan::try_from(object.plan).expect("valid plan"),
1605            })
1606            .collect();
1607
1608        let augmented_dataflow = DataflowDescription {
1609            source_imports,
1610            sink_exports,
1611            objects_to_build,
1612            // The rest of the fields are identical
1613            index_imports: dataflow.index_imports,
1614            index_exports: dataflow.index_exports,
1615            as_of: dataflow.as_of.clone(),
1616            until: dataflow.until,
1617            initial_storage_as_of: dataflow.initial_storage_as_of,
1618            refresh_schedule: dataflow.refresh_schedule,
1619            debug_name: dataflow.debug_name,
1620            time_dependence: dataflow.time_dependence,
1621        };
1622
1623        if augmented_dataflow.is_transient() {
1624            tracing::debug!(
1625                name = %augmented_dataflow.debug_name,
1626                import_ids = %augmented_dataflow.display_import_ids(),
1627                export_ids = %augmented_dataflow.display_export_ids(),
1628                as_of = ?augmented_dataflow.as_of.as_ref().unwrap().elements(),
1629                until = ?augmented_dataflow.until.elements(),
1630                "creating dataflow",
1631            );
1632        } else {
1633            tracing::info!(
1634                name = %augmented_dataflow.debug_name,
1635                import_ids = %augmented_dataflow.display_import_ids(),
1636                export_ids = %augmented_dataflow.display_export_ids(),
1637                as_of = ?augmented_dataflow.as_of.as_ref().unwrap().elements(),
1638                until = ?augmented_dataflow.until.elements(),
1639                "creating dataflow",
1640            );
1641        }
1642
1643        // Skip the actual dataflow creation for an empty `as_of`. (Happens e.g. for the
1644        // bootstrapping of a REFRESH AT mat view that is past its last refresh.)
1645        if as_of.is_empty() {
1646            tracing::info!(
1647                name = %augmented_dataflow.debug_name,
1648                "not sending `CreateDataflow`, because of empty `as_of`",
1649            );
1650        } else {
1651            let collections: Vec<_> = augmented_dataflow.export_ids().collect();
1652            self.send(ComputeCommand::CreateDataflow(Box::new(augmented_dataflow)));
1653
1654            for id in collections {
1655                self.maybe_schedule_collection(id);
1656            }
1657        }
1658
1659        Ok(())
1660    }
1661
1662    /// Schedule the identified collection if all its inputs are available.
1663    ///
1664    /// # Panics
1665    ///
1666    /// Panics if the identified collection does not exist.
1667    fn maybe_schedule_collection(&mut self, id: GlobalId) {
1668        let collection = self.expect_collection(id);
1669
1670        // Don't schedule collections twice.
1671        if collection.scheduled {
1672            return;
1673        }
1674
1675        let as_of = collection.read_frontier();
1676
1677        // If the collection has an empty `as_of`, it was either never installed on the replica or
1678        // has since been dropped. In either case the replica does not expect any commands for it.
1679        if as_of.is_empty() {
1680            return;
1681        }
1682
1683        let ready = if id.is_transient() {
1684            // Always schedule transient collections immediately. The assumption is that those are
1685            // created by interactive user commands and we want to schedule them as quickly as
1686            // possible. Inputs might not yet be available, but when they become available, we
1687            // don't need to wait for the controller to become aware and for the scheduling check
1688            // to run again.
1689            true
1690        } else {
1691            // Ignore self-dependencies. Any self-dependencies do not need to be
1692            // available at the as_of for the dataflow to make progress, so we
1693            // can ignore them here. At the moment, only continual tasks have
1694            // self-dependencies, but this logic is correct for any dataflow, so
1695            // we don't special case it to CTs.
1696            let not_self_dep = |x: &GlobalId| *x != id;
1697
1698            // Make sure we never schedule a collection before its input compute collections have
1699            // been scheduled. Scheduling in the wrong order can lead to deadlocks.
1700            let mut deps_scheduled = true;
1701
1702            // Check dependency frontiers to determine if all inputs are
1703            // available. An input is available when its frontier is greater
1704            // than the `as_of`, i.e., all input data up to and including the
1705            // `as_of` has been sealed.
1706            let compute_deps = collection.compute_dependency_ids().filter(not_self_dep);
1707            let mut compute_frontiers = Vec::new();
1708            for id in compute_deps {
1709                let dep = &self.expect_collection(id);
1710                deps_scheduled &= dep.scheduled;
1711                compute_frontiers.push(dep.write_frontier());
1712            }
1713
1714            let storage_deps = collection.storage_dependency_ids().filter(not_self_dep);
1715            let storage_frontiers = self
1716                .storage_collections
1717                .collections_frontiers(storage_deps.collect())
1718                .expect("must exist");
1719            let storage_frontiers = storage_frontiers.into_iter().map(|f| f.write_frontier);
1720
1721            let mut frontiers = compute_frontiers.into_iter().chain(storage_frontiers);
1722            let frontiers_ready =
1723                frontiers.all(|frontier| PartialOrder::less_than(&as_of, &frontier));
1724
1725            deps_scheduled && frontiers_ready
1726        };
1727
1728        if ready {
1729            self.send(ComputeCommand::Schedule(id));
1730            let collection = self.expect_collection_mut(id);
1731            collection.scheduled = true;
1732        }
1733    }
1734
1735    /// Schedule any unscheduled collections that are ready.
1736    fn schedule_collections(&mut self) {
1737        let ids: Vec<_> = self.collections.keys().copied().collect();
1738        for id in ids {
1739            self.maybe_schedule_collection(id);
1740        }
1741    }
1742
1743    /// Drops the read capability for the given collections and allows their resources to be
1744    /// reclaimed.
1745    #[mz_ore::instrument(level = "debug")]
1746    pub fn drop_collections(&mut self, ids: Vec<GlobalId>) -> Result<(), CollectionMissing> {
1747        for id in &ids {
1748            let collection = self.collection_mut(*id)?;
1749
1750            // Mark the collection as dropped to allow it to be removed from the controller state.
1751            collection.dropped = true;
1752
1753            // Drop the implied and warmup read holds to announce that clients are not
1754            // interested in the collection anymore.
1755            collection.implied_read_hold.release();
1756            collection.warmup_read_hold.release();
1757
1758            // If the collection is a subscribe, stop tracking it. This ensures that the controller
1759            // ceases to produce `SubscribeResponse`s for this subscribe.
1760            self.subscribes.remove(id);
1761            // If the collection is a copy to, stop tracking it. This ensures that the controller
1762            // ceases to produce `CopyToResponse`s` for this copy to.
1763            self.copy_tos.remove(id);
1764        }
1765
1766        Ok(())
1767    }
1768
1769    /// Initiate a peek request for the contents of `id` at `timestamp`.
1770    ///
1771    /// If this returns an error, then it didn't modify any `Instance` state.
1772    #[mz_ore::instrument(level = "debug")]
1773    pub fn peek(
1774        &mut self,
1775        peek_target: PeekTarget,
1776        literal_constraints: Option<Vec<Row>>,
1777        uuid: Uuid,
1778        timestamp: Timestamp,
1779        result_desc: RelationDesc,
1780        finishing: RowSetFinishing,
1781        map_filter_project: mz_expr::SafeMfpPlan,
1782        mut read_hold: ReadHold,
1783        target_replica: Option<ReplicaId>,
1784        peek_response_tx: oneshot::Sender<PeekResponse>,
1785    ) -> Result<(), PeekError> {
1786        use PeekError::*;
1787
1788        let target_id = peek_target.id();
1789
1790        // Downgrade the provided read hold to the peek time.
1791        if read_hold.id() != target_id {
1792            return Err(ReadHoldIdMismatch(read_hold.id()));
1793        }
1794        read_hold
1795            .try_downgrade(Antichain::from_elem(timestamp.clone()))
1796            .map_err(|_| ReadHoldInsufficient(target_id))?;
1797
1798        if let Some(target) = target_replica {
1799            if !self.replica_exists(target) {
1800                return Err(ReplicaMissing(target));
1801            }
1802        }
1803
1804        let otel_ctx = OpenTelemetryContext::obtain();
1805
1806        self.peeks.insert(
1807            uuid,
1808            PendingPeek {
1809                target_replica,
1810                // TODO(guswynn): can we just hold the `tracing::Span` here instead?
1811                otel_ctx: otel_ctx.clone(),
1812                requested_at: Instant::now(),
1813                read_hold,
1814                peek_response_tx,
1815                limit: finishing.limit.map(usize::cast_from),
1816                offset: finishing.offset,
1817            },
1818        );
1819
1820        let peek = Peek {
1821            literal_constraints,
1822            uuid,
1823            timestamp,
1824            finishing,
1825            map_filter_project,
1826            // Obtain an `OpenTelemetryContext` from the thread-local tracing
1827            // tree to forward it on to the compute worker.
1828            otel_ctx,
1829            target: peek_target,
1830            result_desc,
1831        };
1832        self.send(ComputeCommand::Peek(Box::new(peek)));
1833
1834        Ok(())
1835    }
1836
1837    /// Cancels an existing peek request.
1838    #[mz_ore::instrument(level = "debug")]
1839    pub fn cancel_peek(&mut self, uuid: Uuid, reason: PeekResponse) {
1840        let Some(peek) = self.peeks.get_mut(&uuid) else {
1841            tracing::warn!("did not find pending peek for {uuid}");
1842            return;
1843        };
1844
1845        let duration = peek.requested_at.elapsed();
1846        self.metrics
1847            .observe_peek_response(&PeekResponse::Canceled, duration);
1848
1849        // Enqueue a notification for the cancellation.
1850        let otel_ctx = peek.otel_ctx.clone();
1851        otel_ctx.attach_as_parent();
1852
1853        self.deliver_response(ComputeControllerResponse::PeekNotification(
1854            uuid,
1855            PeekNotification::Canceled,
1856            otel_ctx,
1857        ));
1858
1859        // Finish the peek.
1860        // This will also propagate the cancellation to the replicas.
1861        self.finish_peek(uuid, reason);
1862    }
1863
1864    /// Assigns a read policy to specific identifiers.
1865    ///
1866    /// The policies are assigned in the order presented, and repeated identifiers should
1867    /// conclude with the last policy. Changing a policy will immediately downgrade the read
1868    /// capability if appropriate, but it will not "recover" the read capability if the prior
1869    /// capability is already ahead of it.
1870    ///
1871    /// Identifiers not present in `policies` retain their existing read policies.
1872    ///
1873    /// It is an error to attempt to set a read policy for a collection that is not readable in the
1874    /// context of compute. At this time, only indexes are readable compute collections.
1875    #[mz_ore::instrument(level = "debug")]
1876    pub fn set_read_policy(
1877        &mut self,
1878        policies: Vec<(GlobalId, ReadPolicy)>,
1879    ) -> Result<(), ReadPolicyError> {
1880        // Do error checking upfront, to avoid introducing inconsistencies between a collection's
1881        // `implied_capability` and `read_capabilities`.
1882        for (id, _policy) in &policies {
1883            let collection = self.collection(*id)?;
1884            if collection.read_policy.is_none() {
1885                return Err(ReadPolicyError::WriteOnlyCollection(*id));
1886            }
1887        }
1888
1889        for (id, new_policy) in policies {
1890            let collection = self.expect_collection_mut(id);
1891            let new_since = new_policy.frontier(collection.write_frontier().borrow());
1892            let _ = collection.implied_read_hold.try_downgrade(new_since);
1893            collection.read_policy = Some(new_policy);
1894        }
1895
1896        Ok(())
1897    }
1898
1899    /// Advance the global write frontier of the given collection.
1900    ///
1901    /// Frontier regressions are gracefully ignored.
1902    ///
1903    /// # Panics
1904    ///
1905    /// Panics if the identified collection does not exist.
1906    #[mz_ore::instrument(level = "debug")]
1907    fn maybe_update_global_write_frontier(
1908        &mut self,
1909        id: GlobalId,
1910        new_frontier: Antichain<Timestamp>,
1911    ) {
1912        let collection = self.expect_collection_mut(id);
1913
1914        let advanced = collection.shared.lock_write_frontier(|f| {
1915            let advanced = PartialOrder::less_than(f, &new_frontier);
1916            if advanced {
1917                f.clone_from(&new_frontier);
1918            }
1919            advanced
1920        });
1921
1922        if !advanced {
1923            return;
1924        }
1925
1926        // Relax the implied read hold according to the read policy.
1927        let new_since = match &collection.read_policy {
1928            Some(read_policy) => {
1929                // For readable collections the read frontier is determined by applying the
1930                // client-provided read policy to the write frontier.
1931                read_policy.frontier(new_frontier.borrow())
1932            }
1933            None => {
1934                // Write-only collections cannot be read within the context of the compute
1935                // controller, so their read frontier only controls the read holds taken on their
1936                // inputs. We can safely downgrade the input read holds to any time less than the
1937                // write frontier.
1938                //
1939                // Note that some write-only collections (continual tasks) need to observe changes
1940                // at their current write frontier during hydration. Thus, we cannot downgrade the
1941                // read frontier to the write frontier and instead step it back by one.
1942                Antichain::from_iter(
1943                    new_frontier
1944                        .iter()
1945                        .map(|t| t.step_back().unwrap_or(Timestamp::MIN)),
1946                )
1947            }
1948        };
1949        let _ = collection.implied_read_hold.try_downgrade(new_since);
1950
1951        // Report the frontier advancement.
1952        self.deliver_response(ComputeControllerResponse::FrontierUpper {
1953            id,
1954            upper: new_frontier,
1955        });
1956    }
1957
1958    /// Apply a collection read hold change.
1959    pub(super) fn apply_read_hold_change(
1960        &mut self,
1961        id: GlobalId,
1962        mut update: ChangeBatch<Timestamp>,
1963    ) {
1964        let Some(collection) = self.collections.get_mut(&id) else {
1965            soft_panic_or_log!(
1966                "read hold change for absent collection (id={id}, changes={update:?})"
1967            );
1968            return;
1969        };
1970
1971        let new_since = collection.shared.lock_read_capabilities(|caps| {
1972            // Sanity check to prevent corrupted `read_capabilities`, which can cause hard-to-debug
1973            // issues (usually stuck read frontiers).
1974            let read_frontier = caps.frontier();
1975            for (time, diff) in update.iter() {
1976                let count = caps.count_for(time) + diff;
1977                assert!(
1978                    count >= 0,
1979                    "invalid read capabilities update: negative capability \
1980             (id={id:?}, read_capabilities={caps:?}, update={update:?})",
1981                );
1982                assert!(
1983                    count == 0 || read_frontier.less_equal(time),
1984                    "invalid read capabilities update: frontier regression \
1985             (id={id:?}, read_capabilities={caps:?}, update={update:?})",
1986                );
1987            }
1988
1989            // Apply read capability updates and learn about resulting changes to the read
1990            // frontier.
1991            let changes = caps.update_iter(update.drain());
1992
1993            let changed = changes.count() > 0;
1994            changed.then(|| caps.frontier().to_owned())
1995        });
1996
1997        let Some(new_since) = new_since else {
1998            return; // read frontier did not change
1999        };
2000
2001        // Propagate read frontier update to dependencies.
2002        for read_hold in collection.compute_dependencies.values_mut() {
2003            read_hold
2004                .try_downgrade(new_since.clone())
2005                .expect("frontiers don't regress");
2006        }
2007        for read_hold in collection.storage_dependencies.values_mut() {
2008            read_hold
2009                .try_downgrade(new_since.clone())
2010                .expect("frontiers don't regress");
2011        }
2012
2013        // Produce `AllowCompaction` command.
2014        self.send(ComputeCommand::AllowCompaction {
2015            id,
2016            frontier: new_since,
2017        });
2018    }
2019
2020    /// Fulfills a registered peek and cleans up associated state.
2021    ///
2022    /// As part of this we:
2023    ///  * Send a `PeekResponse` through the peek's response channel.
2024    ///  * Emit a `CancelPeek` command to instruct replicas to stop spending resources on this
2025    ///    peek, and to allow the `ComputeCommandHistory` to reduce away the corresponding `Peek`
2026    ///    command.
2027    ///  * Remove the read hold for this peek, unblocking compaction that might have waited on it.
2028    fn finish_peek(&mut self, uuid: Uuid, response: PeekResponse) {
2029        let Some(peek) = self.peeks.remove(&uuid) else {
2030            return;
2031        };
2032
2033        // The recipient might not be interested in the peek response anymore, which is fine.
2034        let _ = peek.peek_response_tx.send(response);
2035
2036        // NOTE: We need to send the `CancelPeek` command _before_ we release the peek's read hold
2037        // (by dropping it), to avoid the edge case that caused database-issues#4812.
2038        self.send(ComputeCommand::CancelPeek { uuid });
2039
2040        drop(peek.read_hold);
2041    }
2042
2043    /// Handles a response from a replica. Replica IDs are re-used across replica restarts, so we
2044    /// use the replica epoch to drop stale responses.
2045    fn handle_response(&mut self, (replica_id, epoch, response): ReplicaResponse) {
2046        // Filter responses from non-existing or stale replicas.
2047        if self
2048            .replicas
2049            .get(&replica_id)
2050            .filter(|replica| replica.epoch == epoch)
2051            .is_none()
2052        {
2053            return;
2054        }
2055
2056        // Invariant: the replica exists and has the expected epoch.
2057
2058        match response {
2059            ComputeResponse::Frontiers(id, frontiers) => {
2060                self.handle_frontiers_response(id, frontiers, replica_id);
2061            }
2062            ComputeResponse::PeekResponse(uuid, peek_response, otel_ctx) => {
2063                self.handle_peek_response(uuid, peek_response, otel_ctx, replica_id);
2064            }
2065            ComputeResponse::CopyToResponse(id, response) => {
2066                self.handle_copy_to_response(id, response, replica_id);
2067            }
2068            ComputeResponse::SubscribeResponse(id, response) => {
2069                self.handle_subscribe_response(id, response, replica_id);
2070            }
2071            ComputeResponse::Status(response) => {
2072                self.handle_status_response(response, replica_id);
2073            }
2074        }
2075    }
2076
2077    /// Handle new frontiers, returning any compute response that needs to
2078    /// be sent to the client.
2079    fn handle_frontiers_response(
2080        &mut self,
2081        id: GlobalId,
2082        frontiers: FrontiersResponse,
2083        replica_id: ReplicaId,
2084    ) {
2085        if !self.collections.contains_key(&id) {
2086            soft_panic_or_log!(
2087                "frontiers update for an unknown collection \
2088                 (id={id}, replica_id={replica_id}, frontiers={frontiers:?})"
2089            );
2090            return;
2091        }
2092        let Some(replica) = self.replicas.get_mut(&replica_id) else {
2093            soft_panic_or_log!(
2094                "frontiers update for an unknown replica \
2095                 (replica_id={replica_id}, frontiers={frontiers:?})"
2096            );
2097            return;
2098        };
2099        let Some(replica_collection) = replica.collections.get_mut(&id) else {
2100            soft_panic_or_log!(
2101                "frontiers update for an unknown replica collection \
2102                 (id={id}, replica_id={replica_id}, frontiers={frontiers:?})"
2103            );
2104            return;
2105        };
2106
2107        if let Some(new_frontier) = frontiers.input_frontier {
2108            replica_collection.update_input_frontier(new_frontier.clone());
2109        }
2110        if let Some(new_frontier) = frontiers.output_frontier {
2111            replica_collection.update_output_frontier(new_frontier.clone());
2112        }
2113        if let Some(new_frontier) = frontiers.write_frontier {
2114            replica_collection.update_write_frontier(new_frontier.clone());
2115            self.maybe_update_global_write_frontier(id, new_frontier);
2116        }
2117    }
2118
2119    #[mz_ore::instrument(level = "debug")]
2120    fn handle_peek_response(
2121        &mut self,
2122        uuid: Uuid,
2123        response: PeekResponse,
2124        otel_ctx: OpenTelemetryContext,
2125        replica_id: ReplicaId,
2126    ) {
2127        otel_ctx.attach_as_parent();
2128
2129        // We might not be tracking this peek anymore, because we have served a response already or
2130        // because it was canceled. If this is the case, we ignore the response.
2131        let Some(peek) = self.peeks.get(&uuid) else {
2132            return;
2133        };
2134
2135        // If the peek is targeting a replica, ignore responses from other replicas.
2136        let target_replica = peek.target_replica.unwrap_or(replica_id);
2137        if target_replica != replica_id {
2138            return;
2139        }
2140
2141        let duration = peek.requested_at.elapsed();
2142        self.metrics.observe_peek_response(&response, duration);
2143
2144        let notification = PeekNotification::new(&response, peek.offset, peek.limit);
2145        // NOTE: We use the `otel_ctx` from the response, not the pending peek, because we
2146        // currently want the parent to be whatever the compute worker did with this peek.
2147        self.deliver_response(ComputeControllerResponse::PeekNotification(
2148            uuid,
2149            notification,
2150            otel_ctx,
2151        ));
2152
2153        self.finish_peek(uuid, response)
2154    }
2155
2156    fn handle_copy_to_response(
2157        &mut self,
2158        sink_id: GlobalId,
2159        response: CopyToResponse,
2160        replica_id: ReplicaId,
2161    ) {
2162        if !self.collections.contains_key(&sink_id) {
2163            soft_panic_or_log!(
2164                "received response for an unknown copy-to \
2165                 (sink_id={sink_id}, replica_id={replica_id})",
2166            );
2167            return;
2168        }
2169        let Some(replica) = self.replicas.get_mut(&replica_id) else {
2170            soft_panic_or_log!("copy-to response for an unknown replica (replica_id={replica_id})");
2171            return;
2172        };
2173        let Some(replica_collection) = replica.collections.get_mut(&sink_id) else {
2174            soft_panic_or_log!(
2175                "copy-to response for an unknown replica collection \
2176                 (sink_id={sink_id}, replica_id={replica_id})"
2177            );
2178            return;
2179        };
2180
2181        // Downgrade the replica frontiers, to enable dropping of input read holds and clean up of
2182        // collection state.
2183        // TODO(database-issues#4701): report copy-to frontiers through `Frontiers` responses
2184        replica_collection.update_write_frontier(Antichain::new());
2185        replica_collection.update_input_frontier(Antichain::new());
2186        replica_collection.update_output_frontier(Antichain::new());
2187
2188        // We might not be tracking this COPY TO because we have already returned a response
2189        // from one of the replicas. In that case, we ignore the response.
2190        if !self.copy_tos.remove(&sink_id) {
2191            return;
2192        }
2193
2194        let result = match response {
2195            CopyToResponse::RowCount(count) => Ok(count),
2196            CopyToResponse::Error(error) => Err(anyhow::anyhow!(error)),
2197            // We should never get here: Replicas only drop copy to collections in response
2198            // to the controller allowing them to do so, and when the controller drops a
2199            // copy to it also removes it from the list of tracked copy_tos (see
2200            // [`Instance::drop_collections`]).
2201            CopyToResponse::Dropped => {
2202                tracing::error!(
2203                    %sink_id, %replica_id,
2204                    "received `Dropped` response for a tracked copy to",
2205                );
2206                return;
2207            }
2208        };
2209
2210        self.deliver_response(ComputeControllerResponse::CopyToResponse(sink_id, result));
2211    }
2212
2213    fn handle_subscribe_response(
2214        &mut self,
2215        subscribe_id: GlobalId,
2216        response: SubscribeResponse,
2217        replica_id: ReplicaId,
2218    ) {
2219        if !self.collections.contains_key(&subscribe_id) {
2220            soft_panic_or_log!(
2221                "received response for an unknown subscribe \
2222                 (subscribe_id={subscribe_id}, replica_id={replica_id})",
2223            );
2224            return;
2225        }
2226        let Some(replica) = self.replicas.get_mut(&replica_id) else {
2227            soft_panic_or_log!(
2228                "subscribe response for an unknown replica (replica_id={replica_id})"
2229            );
2230            return;
2231        };
2232        let Some(replica_collection) = replica.collections.get_mut(&subscribe_id) else {
2233            soft_panic_or_log!(
2234                "subscribe response for an unknown replica collection \
2235                 (subscribe_id={subscribe_id}, replica_id={replica_id})"
2236            );
2237            return;
2238        };
2239
2240        // Always apply replica write frontier updates. Even if the subscribe is not tracked
2241        // anymore, there might still be replicas reading from its inputs, so we need to track the
2242        // frontiers until all replicas have advanced to the empty one.
2243        let write_frontier = match &response {
2244            SubscribeResponse::Batch(batch) => batch.upper.clone(),
2245            SubscribeResponse::DroppedAt(_) => Antichain::new(),
2246        };
2247
2248        // For subscribes we downgrade all replica frontiers based on write frontiers. This should
2249        // be fine because the input and output frontier of a subscribe track its write frontier.
2250        // TODO(database-issues#4701): report subscribe frontiers through `Frontiers` responses
2251        replica_collection.update_write_frontier(write_frontier.clone());
2252        replica_collection.update_input_frontier(write_frontier.clone());
2253        replica_collection.update_output_frontier(write_frontier.clone());
2254
2255        // If the subscribe is not tracked, or targets a different replica, there is nothing to do.
2256        let Some(mut subscribe) = self.subscribes.get(&subscribe_id).cloned() else {
2257            return;
2258        };
2259
2260        // Apply a global frontier update.
2261        // If this is a replica-targeted subscribe, it is important that we advance the global
2262        // frontier only based on responses from the targeted replica. Otherwise, another replica
2263        // could advance to the empty frontier, making us drop the subscribe on the targeted
2264        // replica prematurely.
2265        self.maybe_update_global_write_frontier(subscribe_id, write_frontier);
2266
2267        match response {
2268            SubscribeResponse::Batch(batch) => {
2269                let upper = batch.upper;
2270                let mut updates = batch.updates;
2271
2272                // If this batch advances the subscribe's frontier, we emit all updates at times
2273                // greater or equal to the last frontier (to avoid emitting duplicate updates).
2274                if PartialOrder::less_than(&subscribe.frontier, &upper) {
2275                    let lower = std::mem::replace(&mut subscribe.frontier, upper.clone());
2276
2277                    if upper.is_empty() {
2278                        // This subscribe cannot produce more data. Stop tracking it.
2279                        self.subscribes.remove(&subscribe_id);
2280                    } else {
2281                        // This subscribe can produce more data. Update our tracking of it.
2282                        self.subscribes.insert(subscribe_id, subscribe);
2283                    }
2284
2285                    if let Ok(updates) = updates.as_mut() {
2286                        updates.retain_mut(|updates| {
2287                            let offset = updates.times().partition_point(|t| {
2288                                // True for times that are strictly less than lower (and should be skipped)
2289                                // and false otherwise.
2290                                !lower.less_equal(t)
2291                            });
2292                            let (_, past_lower) = std::mem::take(updates).split_at(offset);
2293                            *updates = past_lower;
2294                            updates.len() > 0
2295                        });
2296                    }
2297                    self.deliver_response(ComputeControllerResponse::SubscribeResponse(
2298                        subscribe_id,
2299                        SubscribeBatch {
2300                            lower,
2301                            upper,
2302                            updates,
2303                        },
2304                    ));
2305                }
2306            }
2307            SubscribeResponse::DroppedAt(frontier) => {
2308                // We should never get here: Replicas only drop subscribe collections in response
2309                // to the controller allowing them to do so, and when the controller drops a
2310                // subscribe it also removes it from the list of tracked subscribes (see
2311                // [`Instance::drop_collections`]).
2312                tracing::error!(
2313                    %subscribe_id,
2314                    %replica_id,
2315                    frontier = ?frontier.elements(),
2316                    "received `DroppedAt` response for a tracked subscribe",
2317                );
2318                self.subscribes.remove(&subscribe_id);
2319            }
2320        }
2321    }
2322
2323    fn handle_status_response(&self, response: StatusResponse, _replica_id: ReplicaId) {
2324        match response {
2325            StatusResponse::Placeholder => {}
2326        }
2327    }
2328
2329    /// Return the write frontiers of the dependencies of the given collection.
2330    fn dependency_write_frontiers<'b>(
2331        &'b self,
2332        collection: &'b CollectionState,
2333    ) -> impl Iterator<Item = Antichain<Timestamp>> + 'b {
2334        let compute_frontiers = collection.compute_dependency_ids().filter_map(|dep_id| {
2335            let collection = self.collections.get(&dep_id);
2336            collection.map(|c| c.write_frontier())
2337        });
2338        let storage_frontiers = collection.storage_dependency_ids().filter_map(|dep_id| {
2339            let frontiers = self.storage_collections.collection_frontiers(dep_id).ok();
2340            frontiers.map(|f| f.write_frontier)
2341        });
2342
2343        compute_frontiers.chain(storage_frontiers)
2344    }
2345
2346    /// Return the write frontiers of transitive storage dependencies of the given collection.
2347    fn transitive_storage_dependency_write_frontiers<'b>(
2348        &'b self,
2349        collection: &'b CollectionState,
2350    ) -> impl Iterator<Item = Antichain<Timestamp>> + 'b {
2351        let mut storage_ids: BTreeSet<_> = collection.storage_dependency_ids().collect();
2352        let mut todo: Vec<_> = collection.compute_dependency_ids().collect();
2353        let mut done = BTreeSet::new();
2354
2355        while let Some(id) = todo.pop() {
2356            if done.contains(&id) {
2357                continue;
2358            }
2359            if let Some(dep) = self.collections.get(&id) {
2360                storage_ids.extend(dep.storage_dependency_ids());
2361                todo.extend(dep.compute_dependency_ids())
2362            }
2363            done.insert(id);
2364        }
2365
2366        let storage_frontiers = storage_ids.into_iter().filter_map(|id| {
2367            let frontiers = self.storage_collections.collection_frontiers(id).ok();
2368            frontiers.map(|f| f.write_frontier)
2369        });
2370
2371        storage_frontiers
2372    }
2373
2374    /// Downgrade the warmup capabilities of collections as much as possible.
2375    ///
2376    /// The only requirement we have for a collection's warmup capability is that it is for a time
2377    /// that is available in all of the collection's inputs. For each input the latest time that is
2378    /// the case for is `write_frontier - 1`. So the farthest we can downgrade a collection's
2379    /// warmup capability is the minimum of `write_frontier - 1` of all its inputs.
2380    ///
2381    /// This method expects to be periodically called as part of instance maintenance work.
2382    /// We would like to instead update the warmup capabilities synchronously in response to
2383    /// frontier updates of dependency collections, but that is not generally possible because we
2384    /// don't learn about frontier updates of storage collections synchronously. We could do
2385    /// synchronous updates for compute dependencies, but we refrain from doing for simplicity.
2386    fn downgrade_warmup_capabilities(&mut self) {
2387        let mut new_capabilities = BTreeMap::new();
2388        for (id, collection) in &self.collections {
2389            // For write-only collections that have advanced to the empty frontier, we can drop the
2390            // warmup capability entirely. There is no reason why we would need to hydrate those
2391            // collections again, so being able to warm them up is not useful.
2392            if collection.read_policy.is_none()
2393                && collection.shared.lock_write_frontier(|f| f.is_empty())
2394            {
2395                new_capabilities.insert(*id, Antichain::new());
2396                continue;
2397            }
2398
2399            let mut new_capability = Antichain::new();
2400            for frontier in self.dependency_write_frontiers(collection) {
2401                for time in frontier {
2402                    new_capability.insert(time.step_back().unwrap_or(time));
2403                }
2404            }
2405
2406            new_capabilities.insert(*id, new_capability);
2407        }
2408
2409        for (id, new_capability) in new_capabilities {
2410            let collection = self.expect_collection_mut(id);
2411            let _ = collection.warmup_read_hold.try_downgrade(new_capability);
2412        }
2413    }
2414
2415    /// Forward the implied capabilities of collections, if possible.
2416    ///
2417    /// The implied capability of a collection controls (a) which times are still readable (for
2418    /// indexes) and (b) with which as-of the collection gets installed on a new replica. We are
2419    /// usually not allowed to advance an implied capability beyond the frontier that follows from
2420    /// the collection's read policy applied to its write frontier:
2421    ///
2422    ///  * For sink collections, some external consumer might rely on seeing all distinct times in
2423    ///    the input reflected in the output. If we'd forward the implied capability of a sink,
2424    ///    we'd risk skipping times in the output across replica restarts.
2425    ///  * For index collections, we might make the index unreadable by advancing its read frontier
2426    ///    beyond its write frontier.
2427    ///
2428    /// There is one case where forwarding an implied capability is fine though: an index installed
2429    /// on a cluster that has no replicas. Such indexes are not readable anyway until a new replica
2430    /// is added, so advancing its read frontier can't make it unreadable. We can thus advance the
2431    /// implied capability as long as we make sure that when a new replica is added, the expected
2432    /// relationship between write frontier, read policy, and implied capability can be restored
2433    /// immediately (modulo computation time).
2434    ///
2435    /// Forwarding implied capabilities is not necessary for the correct functioning of the
2436    /// controller but an optimization that is beneficial in two ways:
2437    ///
2438    ///  * It relaxes read holds on inputs to forwarded collections, allowing their compaction.
2439    ///  * It reduces the amount of historical detail new replicas need to process when computing
2440    ///    forwarded collections, as forwarding the implied capability also forwards the corresponding
2441    ///    dataflow as-of.
2442    fn forward_implied_capabilities(&mut self) {
2443        if !ENABLE_PAUSED_CLUSTER_READHOLD_DOWNGRADE.get(&self.dyncfg) {
2444            return;
2445        }
2446        if !self.replicas.is_empty() {
2447            return;
2448        }
2449
2450        let mut new_capabilities = BTreeMap::new();
2451        for (id, collection) in &self.collections {
2452            let Some(read_policy) = &collection.read_policy else {
2453                // Collection is write-only, i.e. a sink.
2454                continue;
2455            };
2456
2457            // When a new replica is started, it will immediately be able to compute all collection
2458            // output up to the write frontier of its transitive storage inputs. So the new implied
2459            // read capability should be the read policy applied to that frontier.
2460            let mut dep_frontier = Antichain::new();
2461            for frontier in self.transitive_storage_dependency_write_frontiers(collection) {
2462                dep_frontier.extend(frontier);
2463            }
2464
2465            let new_capability = read_policy.frontier(dep_frontier.borrow());
2466            if PartialOrder::less_than(collection.implied_read_hold.since(), &new_capability) {
2467                new_capabilities.insert(*id, new_capability);
2468            }
2469        }
2470
2471        for (id, new_capability) in new_capabilities {
2472            let collection = self.expect_collection_mut(id);
2473            let _ = collection.implied_read_hold.try_downgrade(new_capability);
2474        }
2475    }
2476
2477    /// Acquires a `ReadHold` for the identified compute collection.
2478    ///
2479    /// This mirrors the logic used by the controller-side `InstanceState::acquire_read_hold`,
2480    /// but executes on the instance task itself.
2481    pub(super) fn acquire_read_hold(&self, id: GlobalId) -> Result<ReadHold, CollectionMissing> {
2482        // Similarly to InstanceState::acquire_read_hold and StorageCollections::acquire_read_holds,
2483        // we acquire read holds at the earliest possible time rather than returning a copy
2484        // of the implied read hold. This is so that dependents can acquire read holds on
2485        // compute dependencies at frontiers that are held back by other read holds the caller
2486        // has previously taken.
2487        let collection = self.collection(id)?;
2488        let since = collection.shared.lock_read_capabilities(|caps| {
2489            let since = caps.frontier().to_owned();
2490            caps.update_iter(since.iter().map(|t| (t.clone(), 1)));
2491            since
2492        });
2493        let hold = ReadHold::new(id, since, Arc::clone(&self.read_hold_tx));
2494        Ok(hold)
2495    }
2496
2497    /// Process pending maintenance work.
2498    ///
2499    /// This method is invoked periodically by the global controller.
2500    /// It is a good place to perform maintenance work that arises from various controller state
2501    /// changes and that cannot conveniently be handled synchronously with those state changes.
2502    #[mz_ore::instrument(level = "debug")]
2503    pub fn maintain(&mut self) {
2504        self.rehydrate_failed_replicas();
2505        self.downgrade_warmup_capabilities();
2506        self.forward_implied_capabilities();
2507        self.schedule_collections();
2508        self.cleanup_collections();
2509        self.update_frontier_introspection();
2510        self.refresh_state_metrics();
2511        self.refresh_wallclock_lag();
2512    }
2513}
2514
2515/// State maintained about individual compute collections.
2516///
2517/// A compute collection is either an index, or a storage sink, or a subscribe, exported by a
2518/// compute dataflow.
2519#[derive(Debug)]
2520struct CollectionState {
2521    /// If set, this collection is only maintained by the specified replica.
2522    target_replica: Option<ReplicaId>,
2523    /// Whether this collection is a log collection.
2524    ///
2525    /// Log collections are special in that they are only maintained by a subset of all replicas.
2526    log_collection: bool,
2527    /// Whether this collection has been dropped by a controller client.
2528    ///
2529    /// The controller is allowed to remove the `CollectionState` for a collection only when
2530    /// `dropped == true`. Otherwise, clients might still expect to be able to query information
2531    /// about this collection.
2532    dropped: bool,
2533    /// Whether this collection has been scheduled, i.e., the controller has sent a `Schedule`
2534    /// command for it.
2535    scheduled: bool,
2536
2537    /// Whether this collection is in read-only mode.
2538    ///
2539    /// When in read-only mode, the dataflow is not allowed to affect external state (largely persist).
2540    read_only: bool,
2541
2542    /// State shared with the `ComputeController`.
2543    shared: SharedCollectionState,
2544
2545    /// A read hold maintaining the implicit capability of the collection.
2546    ///
2547    /// This capability is kept to ensure that the collection remains readable according to its
2548    /// `read_policy`. It also ensures that read holds on the collection's dependencies are kept at
2549    /// some time not greater than the collection's `write_frontier`, guaranteeing that the
2550    /// collection's next outputs can always be computed without skipping times.
2551    implied_read_hold: ReadHold,
2552    /// A read hold held to enable dataflow warmup.
2553    ///
2554    /// Dataflow warmup is an optimization that allows dataflows to immediately start hydrating
2555    /// even when their next output time (as implied by the `write_frontier`) is in the future.
2556    /// By installing a read capability derived from the write frontiers of the collection's
2557    /// inputs, we ensure that the as-of of new dataflows installed for the collection is at a time
2558    /// that is immediately available, so hydration can begin immediately too.
2559    warmup_read_hold: ReadHold,
2560    /// The policy to use to downgrade `self.implied_read_hold`.
2561    ///
2562    /// If `None`, the collection is a write-only collection (i.e. a sink). For write-only
2563    /// collections, the `implied_read_hold` is only required for maintaining read holds on the
2564    /// inputs, so we can immediately downgrade it to the `write_frontier`.
2565    read_policy: Option<ReadPolicy>,
2566
2567    /// Storage identifiers on which this collection depends, and read holds this collection
2568    /// requires on them.
2569    storage_dependencies: BTreeMap<GlobalId, ReadHold>,
2570    /// Compute identifiers on which this collection depends, and read holds this collection
2571    /// requires on them.
2572    compute_dependencies: BTreeMap<GlobalId, ReadHold>,
2573
2574    /// Introspection state associated with this collection.
2575    introspection: CollectionIntrospection,
2576
2577    /// Frontier wallclock lag measurements stashed until the next `WallclockLagHistogram`
2578    /// introspection update.
2579    ///
2580    /// Keys are `(period, lag, labels)` triples, values are counts.
2581    ///
2582    /// If this is `None`, wallclock lag is not tracked for this collection.
2583    wallclock_lag_histogram_stash: Option<
2584        BTreeMap<
2585            (
2586                WallclockLagHistogramPeriod,
2587                WallclockLag,
2588                BTreeMap<&'static str, String>,
2589            ),
2590            Diff,
2591        >,
2592    >,
2593}
2594
2595impl CollectionState {
2596    /// Creates a new collection state, with an initial read policy valid from `since`.
2597    fn new(
2598        collection_id: GlobalId,
2599        as_of: Antichain<Timestamp>,
2600        shared: SharedCollectionState,
2601        storage_dependencies: BTreeMap<GlobalId, ReadHold>,
2602        compute_dependencies: BTreeMap<GlobalId, ReadHold>,
2603        read_hold_tx: read_holds::ChangeTx,
2604        introspection: CollectionIntrospection,
2605    ) -> Self {
2606        // A collection is not readable before the `as_of`.
2607        let since = as_of.clone();
2608        // A collection won't produce updates for times before the `as_of`.
2609        let upper = as_of;
2610
2611        // Ensure that the provided `shared` is valid for the given `as_of`.
2612        assert!(shared.lock_read_capabilities(|c| c.frontier() == since.borrow()));
2613        assert!(shared.lock_write_frontier(|f| f == &upper));
2614
2615        // Initialize collection read holds.
2616        // Note that the implied read hold was already added to the `read_capabilities` when
2617        // `shared` was created, so we only need to add the warmup read hold here.
2618        let implied_read_hold =
2619            ReadHold::new(collection_id, since.clone(), Arc::clone(&read_hold_tx));
2620        let warmup_read_hold = ReadHold::new(collection_id, since.clone(), read_hold_tx);
2621
2622        let updates = warmup_read_hold.since().iter().map(|t| (t.clone(), 1));
2623        shared.lock_read_capabilities(|c| {
2624            c.update_iter(updates);
2625        });
2626
2627        // In an effort to keep the produced wallclock lag introspection data small and
2628        // predictable, we disable wallclock lag tracking for transient collections, i.e. slow-path
2629        // select indexes and subscribes.
2630        let wallclock_lag_histogram_stash = match collection_id.is_transient() {
2631            true => None,
2632            false => Some(Default::default()),
2633        };
2634
2635        Self {
2636            target_replica: None,
2637            log_collection: false,
2638            dropped: false,
2639            scheduled: false,
2640            read_only: true,
2641            shared,
2642            implied_read_hold,
2643            warmup_read_hold,
2644            read_policy: Some(ReadPolicy::ValidFrom(since)),
2645            storage_dependencies,
2646            compute_dependencies,
2647            introspection,
2648            wallclock_lag_histogram_stash,
2649        }
2650    }
2651
2652    /// Creates a new collection state for a log collection.
2653    fn new_log_collection(
2654        id: GlobalId,
2655        shared: SharedCollectionState,
2656        read_hold_tx: read_holds::ChangeTx,
2657        introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
2658    ) -> Self {
2659        let since = Antichain::from_elem(Timestamp::MIN);
2660        let introspection = CollectionIntrospection::new(
2661            id,
2662            introspection_tx,
2663            since.clone(),
2664            false,
2665            None,
2666            None,
2667            Vec::new(),
2668        );
2669        let mut state = Self::new(
2670            id,
2671            since,
2672            shared,
2673            Default::default(),
2674            Default::default(),
2675            read_hold_tx,
2676            introspection,
2677        );
2678        state.log_collection = true;
2679        // Log collections are created and scheduled implicitly as part of replica initialization.
2680        state.scheduled = true;
2681        state
2682    }
2683
2684    /// Reports the current read frontier.
2685    fn read_frontier(&self) -> Antichain<Timestamp> {
2686        self.shared
2687            .lock_read_capabilities(|c| c.frontier().to_owned())
2688    }
2689
2690    /// Reports the current write frontier.
2691    fn write_frontier(&self) -> Antichain<Timestamp> {
2692        self.shared.lock_write_frontier(|f| f.clone())
2693    }
2694
2695    fn storage_dependency_ids(&self) -> impl Iterator<Item = GlobalId> + '_ {
2696        self.storage_dependencies.keys().copied()
2697    }
2698
2699    fn compute_dependency_ids(&self) -> impl Iterator<Item = GlobalId> + '_ {
2700        self.compute_dependencies.keys().copied()
2701    }
2702}
2703
2704/// Collection state shared with the `ComputeController`.
2705///
2706/// Having this allows certain controller APIs, such as `ComputeController::collection_frontiers`
2707/// and `ComputeController::acquire_read_hold` to be non-`async`. This comes at the cost of
2708/// complexity (by introducing shared mutable state) and performance (by introducing locking). We
2709/// should aim to reduce the amount of shared state over time, rather than expand it.
2710///
2711/// Note that [`SharedCollectionState`]s are initialized by the `ComputeController` prior to the
2712/// collection's creation in the [`Instance`]. This is to allow compute clients to query frontiers
2713/// and take new read holds immediately, without having to wait for the [`Instance`] to update.
2714#[derive(Clone, Debug)]
2715pub(super) struct SharedCollectionState {
2716    /// Accumulation of read capabilities for the collection.
2717    ///
2718    /// This accumulation contains the capabilities held by all [`ReadHold`]s given out for the
2719    /// collection, including `implied_read_hold` and `warmup_read_hold`.
2720    ///
2721    /// NOTE: This field may only be modified by [`Instance::apply_read_hold_change`],
2722    /// [`Instance::acquire_read_hold`], and `ComputeController::acquire_read_hold`.
2723    /// Nobody else should modify read capabilities directly. Instead, collection users should
2724    /// manage read holds through [`ReadHold`] objects acquired through
2725    /// `ComputeController::acquire_read_hold`.
2726    ///
2727    /// TODO(teskje): Restructure the code to enforce the above in the type system.
2728    read_capabilities: Arc<Mutex<MutableAntichain<Timestamp>>>,
2729    /// The write frontier of this collection.
2730    write_frontier: Arc<Mutex<Antichain<Timestamp>>>,
2731}
2732
2733impl SharedCollectionState {
2734    pub fn new(as_of: Antichain<Timestamp>) -> Self {
2735        // A collection is not readable before the `as_of`.
2736        let since = as_of.clone();
2737        // A collection won't produce updates for times before the `as_of`.
2738        let upper = as_of;
2739
2740        // Initialize read capabilities to the `since`.
2741        // The is the implied read capability. The corresponding [`ReadHold`] is created in
2742        // [`CollectionState::new`].
2743        let mut read_capabilities = MutableAntichain::new();
2744        read_capabilities.update_iter(since.iter().map(|time| (time.clone(), 1)));
2745
2746        Self {
2747            read_capabilities: Arc::new(Mutex::new(read_capabilities)),
2748            write_frontier: Arc::new(Mutex::new(upper)),
2749        }
2750    }
2751
2752    pub fn lock_read_capabilities<F, R>(&self, f: F) -> R
2753    where
2754        F: FnOnce(&mut MutableAntichain<Timestamp>) -> R,
2755    {
2756        let mut caps = self.read_capabilities.lock().expect("poisoned");
2757        f(&mut *caps)
2758    }
2759
2760    pub fn lock_write_frontier<F, R>(&self, f: F) -> R
2761    where
2762        F: FnOnce(&mut Antichain<Timestamp>) -> R,
2763    {
2764        let mut frontier = self.write_frontier.lock().expect("poisoned");
2765        f(&mut *frontier)
2766    }
2767}
2768
2769/// Manages certain introspection relations associated with a collection. Upon creation, it adds
2770/// rows to introspection relations. When dropped, it retracts its managed rows.
2771#[derive(Debug)]
2772struct CollectionIntrospection {
2773    /// The ID of the compute collection.
2774    collection_id: GlobalId,
2775    /// A channel through which introspection updates are delivered.
2776    introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
2777    /// Introspection state for `IntrospectionType::Frontiers`.
2778    ///
2779    /// `Some` if the collection does _not_ sink into a storage collection (i.e. is not an MV). If
2780    /// the collection sinks into storage, the storage controller reports its frontiers instead.
2781    frontiers: Option<FrontiersIntrospectionState>,
2782    /// Introspection state for `IntrospectionType::ComputeMaterializedViewRefreshes`.
2783    ///
2784    /// `Some` if the collection is a REFRESH MV.
2785    refresh: Option<RefreshIntrospectionState>,
2786    /// The IDs of the collection's dependencies, for `IntrospectionType::ComputeDependencies`.
2787    dependency_ids: Vec<GlobalId>,
2788}
2789
2790impl CollectionIntrospection {
2791    fn new(
2792        collection_id: GlobalId,
2793        introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
2794        as_of: Antichain<Timestamp>,
2795        storage_sink: bool,
2796        initial_as_of: Option<Antichain<Timestamp>>,
2797        refresh_schedule: Option<RefreshSchedule>,
2798        dependency_ids: Vec<GlobalId>,
2799    ) -> Self {
2800        let refresh =
2801            match (refresh_schedule, initial_as_of) {
2802                (Some(refresh_schedule), Some(initial_as_of)) => Some(
2803                    RefreshIntrospectionState::new(refresh_schedule, initial_as_of, &as_of),
2804                ),
2805                (refresh_schedule, _) => {
2806                    // If we have a `refresh_schedule`, then the collection is a MV, so we should also have
2807                    // an `initial_as_of`.
2808                    soft_assert_or_log!(
2809                        refresh_schedule.is_none(),
2810                        "`refresh_schedule` without an `initial_as_of`: {collection_id}"
2811                    );
2812                    None
2813                }
2814            };
2815        let frontiers = (!storage_sink).then(|| FrontiersIntrospectionState::new(as_of));
2816
2817        let self_ = Self {
2818            collection_id,
2819            introspection_tx,
2820            frontiers,
2821            refresh,
2822            dependency_ids,
2823        };
2824
2825        self_.report_initial_state();
2826        self_
2827    }
2828
2829    /// Reports the initial introspection state.
2830    fn report_initial_state(&self) {
2831        if let Some(frontiers) = &self.frontiers {
2832            let row = frontiers.row_for_collection(self.collection_id);
2833            let updates = vec![(row, Diff::ONE)];
2834            self.send(IntrospectionType::Frontiers, updates);
2835        }
2836
2837        if let Some(refresh) = &self.refresh {
2838            let row = refresh.row_for_collection(self.collection_id);
2839            let updates = vec![(row, Diff::ONE)];
2840            self.send(IntrospectionType::ComputeMaterializedViewRefreshes, updates);
2841        }
2842
2843        if !self.dependency_ids.is_empty() {
2844            let updates = self.dependency_rows(Diff::ONE);
2845            self.send(IntrospectionType::ComputeDependencies, updates);
2846        }
2847    }
2848
2849    /// Produces rows for the `ComputeDependencies` introspection relation.
2850    fn dependency_rows(&self, diff: Diff) -> Vec<(Row, Diff)> {
2851        self.dependency_ids
2852            .iter()
2853            .map(|dependency_id| {
2854                let row = Row::pack_slice(&[
2855                    Datum::String(&self.collection_id.to_string()),
2856                    Datum::String(&dependency_id.to_string()),
2857                ]);
2858                (row, diff)
2859            })
2860            .collect()
2861    }
2862
2863    /// Observe the given current collection frontiers and update the introspection state as
2864    /// necessary.
2865    fn observe_frontiers(
2866        &mut self,
2867        read_frontier: &Antichain<Timestamp>,
2868        write_frontier: &Antichain<Timestamp>,
2869    ) {
2870        self.update_frontier_introspection(read_frontier, write_frontier);
2871        self.update_refresh_introspection(write_frontier);
2872    }
2873
2874    fn update_frontier_introspection(
2875        &mut self,
2876        read_frontier: &Antichain<Timestamp>,
2877        write_frontier: &Antichain<Timestamp>,
2878    ) {
2879        let Some(frontiers) = &mut self.frontiers else {
2880            return;
2881        };
2882
2883        if &frontiers.read_frontier == read_frontier && &frontiers.write_frontier == write_frontier
2884        {
2885            return; // no change
2886        };
2887
2888        let retraction = frontiers.row_for_collection(self.collection_id);
2889        frontiers.update(read_frontier, write_frontier);
2890        let insertion = frontiers.row_for_collection(self.collection_id);
2891        let updates = vec![(retraction, Diff::MINUS_ONE), (insertion, Diff::ONE)];
2892        self.send(IntrospectionType::Frontiers, updates);
2893    }
2894
2895    fn update_refresh_introspection(&mut self, write_frontier: &Antichain<Timestamp>) {
2896        let Some(refresh) = &mut self.refresh else {
2897            return;
2898        };
2899
2900        let retraction = refresh.row_for_collection(self.collection_id);
2901        refresh.frontier_update(write_frontier);
2902        let insertion = refresh.row_for_collection(self.collection_id);
2903
2904        if retraction == insertion {
2905            return; // no change
2906        }
2907
2908        let updates = vec![(retraction, Diff::MINUS_ONE), (insertion, Diff::ONE)];
2909        self.send(IntrospectionType::ComputeMaterializedViewRefreshes, updates);
2910    }
2911
2912    fn send(&self, introspection_type: IntrospectionType, updates: Vec<(Row, Diff)>) {
2913        // Failure to send means the `ComputeController` has been dropped and doesn't care about
2914        // introspection updates anymore.
2915        let _ = self.introspection_tx.send((introspection_type, updates));
2916    }
2917}
2918
2919impl Drop for CollectionIntrospection {
2920    fn drop(&mut self) {
2921        // Retract collection frontiers.
2922        if let Some(frontiers) = &self.frontiers {
2923            let row = frontiers.row_for_collection(self.collection_id);
2924            let updates = vec![(row, Diff::MINUS_ONE)];
2925            self.send(IntrospectionType::Frontiers, updates);
2926        }
2927
2928        // Retract MV refresh state.
2929        if let Some(refresh) = &self.refresh {
2930            let retraction = refresh.row_for_collection(self.collection_id);
2931            let updates = vec![(retraction, Diff::MINUS_ONE)];
2932            self.send(IntrospectionType::ComputeMaterializedViewRefreshes, updates);
2933        }
2934
2935        // Retract collection dependencies.
2936        if !self.dependency_ids.is_empty() {
2937            let updates = self.dependency_rows(Diff::MINUS_ONE);
2938            self.send(IntrospectionType::ComputeDependencies, updates);
2939        }
2940    }
2941}
2942
2943#[derive(Debug)]
2944struct FrontiersIntrospectionState {
2945    read_frontier: Antichain<Timestamp>,
2946    write_frontier: Antichain<Timestamp>,
2947}
2948
2949impl FrontiersIntrospectionState {
2950    fn new(as_of: Antichain<Timestamp>) -> Self {
2951        Self {
2952            read_frontier: as_of.clone(),
2953            write_frontier: as_of,
2954        }
2955    }
2956
2957    /// Return a `Row` reflecting the current collection frontiers.
2958    fn row_for_collection(&self, collection_id: GlobalId) -> Row {
2959        let read_frontier = self
2960            .read_frontier
2961            .as_option()
2962            .map_or(Datum::Null, |ts| ts.clone().into());
2963        let write_frontier = self
2964            .write_frontier
2965            .as_option()
2966            .map_or(Datum::Null, |ts| ts.clone().into());
2967        Row::pack_slice(&[
2968            Datum::String(&collection_id.to_string()),
2969            read_frontier,
2970            write_frontier,
2971        ])
2972    }
2973
2974    /// Update the introspection state with the given new frontiers.
2975    fn update(
2976        &mut self,
2977        read_frontier: &Antichain<Timestamp>,
2978        write_frontier: &Antichain<Timestamp>,
2979    ) {
2980        if read_frontier != &self.read_frontier {
2981            self.read_frontier.clone_from(read_frontier);
2982        }
2983        if write_frontier != &self.write_frontier {
2984            self.write_frontier.clone_from(write_frontier);
2985        }
2986    }
2987}
2988
2989/// Information needed to compute introspection updates for a REFRESH materialized view when the
2990/// write frontier advances.
2991#[derive(Debug)]
2992struct RefreshIntrospectionState {
2993    // Immutable properties of the MV
2994    refresh_schedule: RefreshSchedule,
2995    initial_as_of: Antichain<Timestamp>,
2996    // Refresh state
2997    next_refresh: Datum<'static>,           // Null or an MzTimestamp
2998    last_completed_refresh: Datum<'static>, // Null or an MzTimestamp
2999}
3000
3001impl RefreshIntrospectionState {
3002    /// Return a `Row` reflecting the current refresh introspection state.
3003    fn row_for_collection(&self, collection_id: GlobalId) -> Row {
3004        Row::pack_slice(&[
3005            Datum::String(&collection_id.to_string()),
3006            self.last_completed_refresh,
3007            self.next_refresh,
3008        ])
3009    }
3010}
3011
3012impl RefreshIntrospectionState {
3013    /// Construct a new [`RefreshIntrospectionState`], and apply an initial `frontier_update()` at
3014    /// the `upper`.
3015    fn new(
3016        refresh_schedule: RefreshSchedule,
3017        initial_as_of: Antichain<Timestamp>,
3018        upper: &Antichain<Timestamp>,
3019    ) -> Self {
3020        let mut self_ = Self {
3021            refresh_schedule: refresh_schedule.clone(),
3022            initial_as_of: initial_as_of.clone(),
3023            next_refresh: Datum::Null,
3024            last_completed_refresh: Datum::Null,
3025        };
3026        self_.frontier_update(upper);
3027        self_
3028    }
3029
3030    /// Should be called whenever the write frontier of the collection advances. It updates the
3031    /// state that should be recorded in introspection relations, but doesn't send the updates yet.
3032    fn frontier_update(&mut self, write_frontier: &Antichain<Timestamp>) {
3033        if write_frontier.is_empty() {
3034            self.last_completed_refresh =
3035                if let Some(last_refresh) = self.refresh_schedule.last_refresh() {
3036                    last_refresh.into()
3037                } else {
3038                    // If there is no last refresh, then we have a `REFRESH EVERY`, in which case
3039                    // the saturating roundup puts a refresh at the maximum possible timestamp.
3040                    Timestamp::MAX.into()
3041                };
3042            self.next_refresh = Datum::Null;
3043        } else {
3044            if PartialOrder::less_equal(write_frontier, &self.initial_as_of) {
3045                // We are before the first refresh.
3046                self.last_completed_refresh = Datum::Null;
3047                let initial_as_of = self.initial_as_of.as_option().expect(
3048                    "initial_as_of can't be [], because then there would be no refreshes at all",
3049                );
3050                let first_refresh = self
3051                    .refresh_schedule
3052                    .round_up_timestamp(*initial_as_of)
3053                    .expect("sequencing makes sure that REFRESH MVs always have a first refresh");
3054                soft_assert_or_log!(
3055                    first_refresh == *initial_as_of,
3056                    "initial_as_of should be set to the first refresh"
3057                );
3058                self.next_refresh = first_refresh.into();
3059            } else {
3060                // The first refresh has already happened.
3061                let write_frontier = write_frontier.as_option().expect("checked above");
3062                self.last_completed_refresh = self
3063                    .refresh_schedule
3064                    .round_down_timestamp_m1(*write_frontier)
3065                    .map_or_else(
3066                        || {
3067                            soft_panic_or_log!(
3068                                "rounding down should have returned the first refresh or later"
3069                            );
3070                            Datum::Null
3071                        },
3072                        |last_completed_refresh| last_completed_refresh.into(),
3073                    );
3074                self.next_refresh = write_frontier.clone().into();
3075            }
3076        }
3077    }
3078}
3079
3080/// A note of an outstanding peek response.
3081#[derive(Debug)]
3082struct PendingPeek {
3083    /// For replica-targeted peeks, this specifies the replica whose response we should pass on.
3084    ///
3085    /// If this value is `None`, we pass on the first response.
3086    target_replica: Option<ReplicaId>,
3087    /// The OpenTelemetry context for this peek.
3088    otel_ctx: OpenTelemetryContext,
3089    /// The time at which the peek was requested.
3090    ///
3091    /// Used to track peek durations.
3092    requested_at: Instant,
3093    /// The read hold installed to serve this peek.
3094    read_hold: ReadHold,
3095    /// The channel to send peek results.
3096    peek_response_tx: oneshot::Sender<PeekResponse>,
3097    /// An optional limit of the peek's result size.
3098    limit: Option<usize>,
3099    /// The offset into the peek's result.
3100    offset: usize,
3101}
3102
3103#[derive(Debug, Clone)]
3104struct ActiveSubscribe {
3105    /// Current upper frontier of this subscribe.
3106    frontier: Antichain<Timestamp>,
3107}
3108
3109impl Default for ActiveSubscribe {
3110    fn default() -> Self {
3111        Self {
3112            frontier: Antichain::from_elem(Timestamp::MIN),
3113        }
3114    }
3115}
3116
3117/// State maintained about individual replicas.
3118#[derive(Debug)]
3119struct ReplicaState {
3120    /// The ID of the replica.
3121    id: ReplicaId,
3122    /// Client for the running replica task.
3123    client: ReplicaClient,
3124    /// The replica configuration.
3125    config: ReplicaConfig,
3126    /// Replica metrics.
3127    metrics: ReplicaMetrics,
3128    /// A channel through which introspection updates are delivered.
3129    introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
3130    /// Per-replica collection state.
3131    collections: BTreeMap<GlobalId, ReplicaCollectionState>,
3132    /// The epoch of the replica.
3133    epoch: u64,
3134}
3135
3136impl ReplicaState {
3137    fn new(
3138        id: ReplicaId,
3139        client: ReplicaClient,
3140        config: ReplicaConfig,
3141        metrics: ReplicaMetrics,
3142        introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
3143        epoch: u64,
3144    ) -> Self {
3145        Self {
3146            id,
3147            client,
3148            config,
3149            metrics,
3150            introspection_tx,
3151            epoch,
3152            collections: Default::default(),
3153        }
3154    }
3155
3156    /// Add a collection to the replica state.
3157    ///
3158    /// # Panics
3159    ///
3160    /// Panics if a collection with the same ID exists already.
3161    fn add_collection(
3162        &mut self,
3163        id: GlobalId,
3164        as_of: Antichain<Timestamp>,
3165        input_read_holds: Vec<ReadHold>,
3166    ) {
3167        let metrics = self.metrics.for_collection(id);
3168        let introspection = ReplicaCollectionIntrospection::new(
3169            self.id,
3170            id,
3171            self.introspection_tx.clone(),
3172            as_of.clone(),
3173        );
3174        let mut state =
3175            ReplicaCollectionState::new(metrics, as_of, introspection, input_read_holds);
3176
3177        // In an effort to keep the produced wallclock lag introspection data small and
3178        // predictable, we disable wallclock lag tracking for transient collections, i.e. slow-path
3179        // select indexes and subscribes.
3180        if id.is_transient() {
3181            state.wallclock_lag_max = None;
3182        }
3183
3184        if let Some(previous) = self.collections.insert(id, state) {
3185            panic!("attempt to add a collection with existing ID {id} (previous={previous:?}");
3186        }
3187    }
3188
3189    /// Remove state for a collection.
3190    fn remove_collection(&mut self, id: GlobalId) -> Option<ReplicaCollectionState> {
3191        self.collections.remove(&id)
3192    }
3193
3194    /// Returns the per-replica state of a collection this replica hosts.
3195    ///
3196    /// # Panics
3197    ///
3198    /// Panics if the replica does not host the collection. Callers obtain the
3199    /// replica from [`Instance::replicas_hosting`], which guarantees it does.
3200    fn expect_collection(&self, id: GlobalId) -> &ReplicaCollectionState {
3201        self.collections
3202            .get(&id)
3203            .expect("hosting replica must have per-replica collection state")
3204    }
3205
3206    /// Returns whether all replica frontiers of the given collection are empty.
3207    fn collection_frontiers_empty(&self, id: GlobalId) -> bool {
3208        self.collections.get(&id).map_or(true, |c| {
3209            c.write_frontier.is_empty()
3210                && c.input_frontier.is_empty()
3211                && c.output_frontier.is_empty()
3212        })
3213    }
3214
3215    /// Returns the state of the [`ReplicaState`] formatted as JSON.
3216    ///
3217    /// The returned value is not guaranteed to be stable and may change at any point in time.
3218    #[mz_ore::instrument(level = "debug")]
3219    pub fn dump(&self) -> Result<serde_json::Value, anyhow::Error> {
3220        // Note: We purposefully use the `Debug` formatting for the value of all fields in the
3221        // returned object as a tradeoff between usability and stability. `serde_json` will fail
3222        // to serialize an object if the keys aren't strings, so `Debug` formatting the values
3223        // prevents a future unrelated change from silently breaking this method.
3224
3225        // Destructure `self` here so we don't forget to consider dumping newly added fields.
3226        let Self {
3227            id,
3228            client: _,
3229            config: _,
3230            metrics: _,
3231            introspection_tx: _,
3232            epoch,
3233            collections,
3234        } = self;
3235
3236        let collections: BTreeMap<_, _> = collections
3237            .iter()
3238            .map(|(id, collection)| (id.to_string(), format!("{collection:?}")))
3239            .collect();
3240
3241        Ok(serde_json::json!({
3242            "id": id.to_string(),
3243            "collections": collections,
3244            "epoch": epoch,
3245        }))
3246    }
3247}
3248
3249#[derive(Debug)]
3250struct ReplicaCollectionState {
3251    /// The replica write frontier of this collection.
3252    ///
3253    /// See [`FrontiersResponse::write_frontier`].
3254    write_frontier: Antichain<Timestamp>,
3255    /// The replica input frontier of this collection.
3256    ///
3257    /// See [`FrontiersResponse::input_frontier`].
3258    input_frontier: Antichain<Timestamp>,
3259    /// The replica output frontier of this collection.
3260    ///
3261    /// See [`FrontiersResponse::output_frontier`].
3262    output_frontier: Antichain<Timestamp>,
3263
3264    /// Metrics tracked for this collection.
3265    ///
3266    /// If this is `None`, no metrics are collected.
3267    metrics: Option<ReplicaCollectionMetrics>,
3268    /// As-of frontier with which this collection was installed on the replica.
3269    as_of: Antichain<Timestamp>,
3270    /// Tracks introspection state for this collection.
3271    introspection: ReplicaCollectionIntrospection,
3272    /// Read holds on storage inputs to this collection.
3273    ///
3274    /// These read holds are kept to ensure that the replica is able to read from storage inputs at
3275    /// all times it hasn't read yet. We only need to install read holds for storage inputs since
3276    /// compaction of compute inputs is implicitly held back by Timely/DD.
3277    input_read_holds: Vec<ReadHold>,
3278
3279    /// Maximum frontier wallclock lag since the last `WallclockLagHistory` introspection update.
3280    ///
3281    /// If this is `None`, wallclock lag is not tracked for this collection.
3282    wallclock_lag_max: Option<WallclockLag>,
3283}
3284
3285impl ReplicaCollectionState {
3286    fn new(
3287        metrics: Option<ReplicaCollectionMetrics>,
3288        as_of: Antichain<Timestamp>,
3289        introspection: ReplicaCollectionIntrospection,
3290        input_read_holds: Vec<ReadHold>,
3291    ) -> Self {
3292        Self {
3293            write_frontier: as_of.clone(),
3294            input_frontier: as_of.clone(),
3295            output_frontier: as_of.clone(),
3296            metrics,
3297            as_of,
3298            introspection,
3299            input_read_holds,
3300            wallclock_lag_max: Some(WallclockLag::MIN),
3301        }
3302    }
3303
3304    /// Returns whether this collection is hydrated.
3305    fn hydrated(&self) -> bool {
3306        // If the observed frontier is greater than the collection's as-of, the collection has
3307        // produced some output and is therefore hydrated.
3308        //
3309        // We need to consider the edge case where the as-of is the empty frontier. Such an as-of
3310        // is not useful for indexes, because they wouldn't be readable. For write-only
3311        // collections, an empty as-of means that the collection has been fully written and no new
3312        // dataflow needs to be created for it. Consequently, no hydration will happen either.
3313        //
3314        // Based on this, we could respond in two ways:
3315        //  * `false`, as in "the dataflow was never created"
3316        //  * `true`, as in "the dataflow completed immediately"
3317        //
3318        // Since hydration is often used as a measure of dataflow progress and we don't want to
3319        // give the impression that certain dataflows are somehow stuck when they are not, we go
3320        // with the second interpretation here.
3321        self.as_of.is_empty() || PartialOrder::less_than(&self.as_of, &self.output_frontier)
3322    }
3323
3324    /// Updates the replica write frontier of this collection.
3325    fn update_write_frontier(&mut self, new_frontier: Antichain<Timestamp>) {
3326        if PartialOrder::less_than(&new_frontier, &self.write_frontier) {
3327            soft_panic_or_log!(
3328                "replica collection write frontier regression (old={:?}, new={new_frontier:?})",
3329                self.write_frontier,
3330            );
3331            return;
3332        } else if new_frontier == self.write_frontier {
3333            return;
3334        }
3335
3336        self.write_frontier = new_frontier;
3337    }
3338
3339    /// Updates the replica input frontier of this collection.
3340    fn update_input_frontier(&mut self, new_frontier: Antichain<Timestamp>) {
3341        if PartialOrder::less_than(&new_frontier, &self.input_frontier) {
3342            soft_panic_or_log!(
3343                "replica collection input frontier regression (old={:?}, new={new_frontier:?})",
3344                self.input_frontier,
3345            );
3346            return;
3347        } else if new_frontier == self.input_frontier {
3348            return;
3349        }
3350
3351        self.input_frontier = new_frontier;
3352
3353        // Relax our read holds on collection inputs.
3354        for read_hold in &mut self.input_read_holds {
3355            let result = read_hold.try_downgrade(self.input_frontier.clone());
3356            soft_assert_or_log!(
3357                result.is_ok(),
3358                "read hold downgrade failed (read_hold={read_hold:?}, new_since={:?})",
3359                self.input_frontier,
3360            );
3361        }
3362    }
3363
3364    /// Updates the replica output frontier of this collection.
3365    fn update_output_frontier(&mut self, new_frontier: Antichain<Timestamp>) {
3366        if PartialOrder::less_than(&new_frontier, &self.output_frontier) {
3367            soft_panic_or_log!(
3368                "replica collection output frontier regression (old={:?}, new={new_frontier:?})",
3369                self.output_frontier,
3370            );
3371            return;
3372        } else if new_frontier == self.output_frontier {
3373            return;
3374        }
3375
3376        self.output_frontier = new_frontier;
3377    }
3378}
3379
3380/// Maintains the introspection state for a given replica and collection, and ensures that reported
3381/// introspection data is retracted when the collection is dropped.
3382#[derive(Debug)]
3383struct ReplicaCollectionIntrospection {
3384    /// The ID of the replica.
3385    replica_id: ReplicaId,
3386    /// The ID of the compute collection.
3387    collection_id: GlobalId,
3388    /// The collection's reported replica write frontier.
3389    write_frontier: Antichain<Timestamp>,
3390    /// A channel through which introspection updates are delivered.
3391    introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
3392}
3393
3394impl ReplicaCollectionIntrospection {
3395    /// Create a new `HydrationState` and initialize introspection.
3396    fn new(
3397        replica_id: ReplicaId,
3398        collection_id: GlobalId,
3399        introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
3400        as_of: Antichain<Timestamp>,
3401    ) -> Self {
3402        let self_ = Self {
3403            replica_id,
3404            collection_id,
3405            write_frontier: as_of,
3406            introspection_tx,
3407        };
3408
3409        self_.report_initial_state();
3410        self_
3411    }
3412
3413    /// Reports the initial introspection state.
3414    fn report_initial_state(&self) {
3415        let row = self.write_frontier_row();
3416        let updates = vec![(row, Diff::ONE)];
3417        self.send(IntrospectionType::ReplicaFrontiers, updates);
3418    }
3419
3420    /// Observe the given current write frontier and update the introspection state as necessary.
3421    fn observe_frontier(&mut self, write_frontier: &Antichain<Timestamp>) {
3422        if self.write_frontier == *write_frontier {
3423            return; // no change
3424        }
3425
3426        let retraction = self.write_frontier_row();
3427        self.write_frontier.clone_from(write_frontier);
3428        let insertion = self.write_frontier_row();
3429
3430        let updates = vec![(retraction, Diff::MINUS_ONE), (insertion, Diff::ONE)];
3431        self.send(IntrospectionType::ReplicaFrontiers, updates);
3432    }
3433
3434    /// Return a `Row` reflecting the current replica write frontier.
3435    fn write_frontier_row(&self) -> Row {
3436        let write_frontier = self
3437            .write_frontier
3438            .as_option()
3439            .map_or(Datum::Null, |ts| ts.clone().into());
3440        Row::pack_slice(&[
3441            Datum::String(&self.collection_id.to_string()),
3442            Datum::String(&self.replica_id.to_string()),
3443            write_frontier,
3444        ])
3445    }
3446
3447    fn send(&self, introspection_type: IntrospectionType, updates: Vec<(Row, Diff)>) {
3448        // Failure to send means the `ComputeController` has been dropped and doesn't care about
3449        // introspection updates anymore.
3450        let _ = self.introspection_tx.send((introspection_type, updates));
3451    }
3452}
3453
3454impl Drop for ReplicaCollectionIntrospection {
3455    fn drop(&mut self) {
3456        // Retract the write frontier.
3457        let row = self.write_frontier_row();
3458        let updates = vec![(row, Diff::MINUS_ONE)];
3459        self.send(IntrospectionType::ReplicaFrontiers, updates);
3460    }
3461}
3462
3463/// Classifies one collection's readiness over every replica hosting it.
3464///
3465/// Hydration and sufficient progress must belong to the same target replica.
3466/// No reference replicas yields `MIN`, while a completed reference yields the
3467/// empty frontier. Those cases must remain distinct.
3468fn classify_collection_readiness<'a, I>(
3469    replicas: I,
3470    target_replica_ids: &BTreeSet<ReplicaId>,
3471    reference_replica_ids: &BTreeSet<ReplicaId>,
3472    allowed_lag: Option<Timestamp>,
3473) -> CollectionReadiness
3474where
3475    I: Iterator<Item = (ReplicaId, &'a ReplicaCollectionState)> + Clone,
3476{
3477    let lag = allowed_lag.map(|allowed_lag| {
3478        let mut reference = Antichain::from_elem(Timestamp::MIN);
3479        for (_, state) in replicas
3480            .clone()
3481            .filter(|(id, _)| reference_replica_ids.contains(id))
3482        {
3483            reference.join_assign(&state.output_frontier);
3484        }
3485        (reference, allowed_lag)
3486    });
3487
3488    let mut result = CollectionReadiness::Unhydrated;
3489    for (_, state) in replicas.filter(|(id, _)| target_replica_ids.contains(id)) {
3490        match CollectionReadiness::classify(
3491            state.hydrated(),
3492            &state.output_frontier,
3493            lag.as_ref().map(|(reference, lag)| (reference, *lag)),
3494        ) {
3495            CollectionReadiness::Ready => return CollectionReadiness::Ready,
3496            CollectionReadiness::Lagging { lag } => {
3497                // Report the closest hydrated target, since any ready target
3498                // suffices. A finite gap is closer than awaiting completion.
3499                let lag = match result {
3500                    CollectionReadiness::Lagging {
3501                        lag: Some(previous),
3502                    } => Some(previous.min(lag.unwrap_or(u64::MAX))),
3503                    _ => lag,
3504                };
3505                result = CollectionReadiness::Lagging { lag };
3506            }
3507            CollectionReadiness::Unhydrated => {}
3508        }
3509    }
3510    result
3511}
3512
3513#[cfg(test)]
3514mod tests {
3515    use std::collections::{BTreeMap, BTreeSet};
3516
3517    use mz_compute_types::dyncfgs::{ENABLE_COLUMN_PAGED_BATCHER, ENABLE_MZ_JOIN_CORE};
3518    use mz_dyncfg::{ConfigSet, ConfigUpdates, ConfigVal};
3519    use mz_persist_types::PersistLocation;
3520    use mz_repr::{GlobalId, Timestamp};
3521    use timely::progress::Antichain;
3522    use tokio::sync::mpsc;
3523
3524    use crate::protocol::command::{ComputeCommand, InstanceConfig};
3525
3526    use super::{
3527        CollectionReadiness, Instance, ReplicaCollectionIntrospection, ReplicaCollectionState,
3528        ReplicaId, classify_collection_readiness,
3529    };
3530
3531    fn ac(ts: u64) -> Antichain<Timestamp> {
3532        Antichain::from_elem(Timestamp::new(ts))
3533    }
3534
3535    fn state(id: ReplicaId, as_of: u64, write: u64, output: u64) -> ReplicaCollectionState {
3536        let (tx, _rx) = mpsc::unbounded_channel();
3537        let introspection =
3538            ReplicaCollectionIntrospection::new(id, GlobalId::User(1), tx, ac(as_of));
3539        let mut state = ReplicaCollectionState::new(None, ac(as_of), introspection, Vec::new());
3540        state.update_write_frontier(ac(write));
3541        state.update_output_frontier(ac(output));
3542        state
3543    }
3544
3545    fn classify(
3546        replicas: &[(ReplicaId, ReplicaCollectionState)],
3547        targets: &[ReplicaId],
3548        references: &[ReplicaId],
3549        lag: Option<Timestamp>,
3550    ) -> CollectionReadiness {
3551        classify_collection_readiness(
3552            replicas.iter().map(|(id, state)| (*id, state)),
3553            &targets.iter().copied().collect::<BTreeSet<_>>(),
3554            &references.iter().copied().collect::<BTreeSet<_>>(),
3555            lag,
3556        )
3557    }
3558
3559    const LAG: Option<Timestamp> = Some(Timestamp::new(60));
3560
3561    /// The regression the lag gate exists for: a pending replica whose dataflow
3562    /// has produced its first output past the as-of, and so reports hydrated,
3563    /// while its output frontier is still far behind the outgoing replica.
3564    #[mz_ore::test]
3565    fn shared_write_does_not_hide_output_lag() {
3566        let reference = ReplicaId::User(1);
3567        let target = ReplicaId::User(2);
3568        let mut replicas = vec![
3569            (reference, state(reference, 100, 10_000, 10_000)),
3570            // The shared MV write frontier is caught up, but this replica has
3571            // only just produced its first output past the as-of.
3572            (target, state(target, 100, 10_000, 101)),
3573        ];
3574        assert_eq!(
3575            classify(&replicas, &[target], &[reference], LAG),
3576            CollectionReadiness::Lagging { lag: Some(9_899) },
3577        );
3578        replicas[1].1.update_output_frontier(ac(9_940));
3579        assert_eq!(
3580            classify(&replicas, &[target], &[reference], LAG),
3581            CollectionReadiness::Ready,
3582        );
3583    }
3584
3585    #[mz_ore::test]
3586    fn refresh_writes_ahead_of_outputs_is_ready() {
3587        let reference = ReplicaId::User(1);
3588        let target = ReplicaId::User(2);
3589        let replicas = vec![
3590            (reference, state(reference, 100, 20_000, 10_000)),
3591            (target, state(target, 100, 20_000, 10_000)),
3592        ];
3593        assert_eq!(
3594            classify(&replicas, &[target], &[reference], LAG),
3595            CollectionReadiness::Ready,
3596        );
3597    }
3598
3599    #[mz_ore::test]
3600    fn furthest_reference_and_bystander_exclusion() {
3601        let trailing = ReplicaId::User(1);
3602        let ahead = ReplicaId::User(2);
3603        let target = ReplicaId::User(3);
3604        let bystander = ReplicaId::User(4);
3605        let replicas = vec![
3606            (trailing, state(trailing, 100, 9_000, 9_000)),
3607            (ahead, state(ahead, 100, 10_000, 10_000)),
3608            (target, state(target, 100, 8_990, 8_990)),
3609            (bystander, state(bystander, 100, 1_000_000, 1_000_000)),
3610        ];
3611        assert_eq!(
3612            classify(&replicas, &[target], &[trailing], LAG),
3613            CollectionReadiness::Ready,
3614        );
3615        assert_eq!(
3616            classify(&replicas, &[target], &[trailing, ahead], LAG),
3617            CollectionReadiness::Lagging { lag: Some(1_010) },
3618        );
3619    }
3620
3621    #[mz_ore::test]
3622    fn one_ready_target_suffices() {
3623        let reference = ReplicaId::User(1);
3624        let unhydrated = ReplicaId::User(2);
3625        let lagging = ReplicaId::User(3);
3626        let ready = ReplicaId::User(4);
3627        let replicas = vec![
3628            (reference, state(reference, 100, 10_000, 10_000)),
3629            // This replica is within the lag allowance but has not advanced
3630            // past its own as-of. It cannot supply progress for the hydrated,
3631            // lagging replica below.
3632            (unhydrated, state(unhydrated, 10_000, 10_000, 10_000)),
3633            (lagging, state(lagging, 100, 1_000, 1_000)),
3634            (ready, state(ready, 100, 9_990, 9_990)),
3635        ];
3636        assert_eq!(
3637            classify(&replicas, &[unhydrated, lagging, ready], &[reference], LAG),
3638            CollectionReadiness::Ready,
3639        );
3640        assert_eq!(
3641            classify(&replicas, &[unhydrated, lagging], &[reference], LAG),
3642            CollectionReadiness::Lagging { lag: Some(9_000) },
3643        );
3644    }
3645
3646    #[mz_ore::test]
3647    fn no_lag_gate_is_hydration_only() {
3648        let reference = ReplicaId::User(1);
3649        let target = ReplicaId::User(2);
3650        let replicas = vec![
3651            (reference, state(reference, 100, 10_000, 10_000)),
3652            (target, state(target, 100, 101, 101)),
3653        ];
3654        assert_eq!(
3655            classify(&replicas, &[target], &[reference], None),
3656            CollectionReadiness::Ready,
3657        );
3658        let unhydrated = vec![(target, state(target, 100, 100, 100))];
3659        assert_eq!(
3660            classify(&unhydrated, &[target], &[], None),
3661            CollectionReadiness::Unhydrated,
3662        );
3663    }
3664
3665    #[mz_ore::test]
3666    fn lag_reports_the_closest_hydrated_target() {
3667        let reference = ReplicaId::User(1);
3668        let far = ReplicaId::User(2);
3669        let close = ReplicaId::User(3);
3670        let mut replicas = vec![
3671            (reference, state(reference, 100, 10_000, 10_000)),
3672            (far, state(far, 100, 1_000, 1_000)),
3673            (close, state(close, 100, 9_000, 9_000)),
3674        ];
3675        for _ in 0..2 {
3676            assert_eq!(
3677                classify(&replicas, &[far, close], &[reference], LAG),
3678                CollectionReadiness::Lagging { lag: Some(1_000) },
3679            );
3680            replicas.reverse();
3681        }
3682    }
3683
3684    #[mz_ore::test]
3685    fn completed_reference_requires_completed_target() {
3686        let reference = ReplicaId::User(1);
3687        let target = ReplicaId::User(2);
3688        let mut replicas = vec![
3689            (reference, state(reference, 100, 1_000, 1_000)),
3690            (target, state(target, 100, 1_000, 1_000)),
3691        ];
3692        replicas[0].1.update_output_frontier(Antichain::new());
3693        assert_eq!(
3694            classify(&replicas, &[target], &[reference], Some(Timestamp::new(0))),
3695            CollectionReadiness::Lagging { lag: None },
3696        );
3697        replicas[1].1.update_output_frontier(Antichain::new());
3698        assert_eq!(
3699            classify(&replicas, &[target], &[reference], Some(Timestamp::new(0))),
3700            CollectionReadiness::Ready,
3701        );
3702    }
3703
3704    #[mz_ore::test]
3705    fn no_reference_means_nothing_to_regress() {
3706        let target = ReplicaId::User(1);
3707        let replicas = vec![(target, state(target, 100, 5_000, 5_000))];
3708        assert_eq!(
3709            classify(&replicas, &[target], &[], Some(Timestamp::new(0))),
3710            CollectionReadiness::Ready,
3711        );
3712        assert_eq!(
3713            classify(&replicas, &[target], &[target], Some(Timestamp::new(0))),
3714            CollectionReadiness::Ready,
3715        );
3716    }
3717
3718    #[mz_ore::test]
3719    fn no_targets_is_unhydrated() {
3720        let reference = ReplicaId::User(1);
3721        let replicas = vec![(reference, state(reference, 100, 10_000, 10_000))];
3722        assert_eq!(
3723            classify(&replicas, &[], &[reference], LAG),
3724            CollectionReadiness::Unhydrated,
3725        );
3726    }
3727
3728    fn create_instance_command() -> ComputeCommand {
3729        ComputeCommand::CreateInstance(Box::new(InstanceConfig {
3730            logging: Default::default(),
3731            expiration_offset: None,
3732            peek_stash_persist_location: PersistLocation::new_in_mem(),
3733            arrangement_dictionary_compression: false,
3734            initial_config: Default::default(),
3735        }))
3736    }
3737
3738    fn initial_config(cmd: &ComputeCommand) -> &ConfigUpdates {
3739        match cmd {
3740            ComputeCommand::CreateInstance(config) => &config.initial_config,
3741            other => panic!("expected CreateInstance, got {other:?}"),
3742        }
3743    }
3744
3745    /// `CreateInstance` is specialized with a full snapshot of the instance-wide dyncfg, so the
3746    /// replica seeds its worker config at create time rather than waiting for the first
3747    /// `UpdateConfiguration`. This is the regression guard for create-time setup observing dyncfg
3748    /// defaults.
3749    #[mz_ore::test]
3750    fn create_instance_snapshots_instance_wide_dyncfg() {
3751        let dyncfg = ConfigSet::default()
3752            .add(&ENABLE_COLUMN_PAGED_BATCHER)
3753            .add(&ENABLE_MZ_JOIN_CORE);
3754        let mut updates = ConfigUpdates::default();
3755        updates.add(&ENABLE_COLUMN_PAGED_BATCHER, true);
3756        updates.add(&ENABLE_MZ_JOIN_CORE, false);
3757        updates.apply(&dyncfg);
3758
3759        // A replica without an override sees exactly the instance-wide values.
3760        let overrides = BTreeMap::new();
3761        let cmd = Instance::specialize_command_for_replica(
3762            create_instance_command(),
3763            ReplicaId::User(1),
3764            &overrides,
3765            &dyncfg,
3766        );
3767        let snapshot = initial_config(&cmd);
3768        assert_eq!(
3769            snapshot.updates.get(ENABLE_COLUMN_PAGED_BATCHER.name()),
3770            Some(&ConfigVal::Bool(true)),
3771        );
3772        assert_eq!(
3773            snapshot.updates.get(ENABLE_MZ_JOIN_CORE.name()),
3774            Some(&ConfigVal::Bool(false)),
3775        );
3776    }
3777
3778    /// A replica-scoped override beats the instance-wide value in the create-time snapshot, so a
3779    /// create-time-frozen scoped flag reaches the replica with its override applied.
3780    #[mz_ore::test]
3781    fn create_instance_snapshot_applies_replica_override() {
3782        let dyncfg = ConfigSet::default().add(&ENABLE_COLUMN_PAGED_BATCHER);
3783        let mut updates = ConfigUpdates::default();
3784        updates.add(&ENABLE_COLUMN_PAGED_BATCHER, true);
3785        updates.apply(&dyncfg);
3786
3787        let replica = ReplicaId::User(1);
3788        let mut override_updates = ConfigUpdates::default();
3789        override_updates.add(&ENABLE_COLUMN_PAGED_BATCHER, false);
3790        let overrides = BTreeMap::from([(replica, override_updates)]);
3791
3792        let cmd = Instance::specialize_command_for_replica(
3793            create_instance_command(),
3794            replica,
3795            &overrides,
3796            &dyncfg,
3797        );
3798        assert_eq!(
3799            initial_config(&cmd)
3800                .updates
3801                .get(ENABLE_COLUMN_PAGED_BATCHER.name()),
3802            Some(&ConfigVal::Bool(false)),
3803            "replica override should win over the instance-wide value",
3804        );
3805    }
3806
3807    /// `UpdateConfiguration` continues to merge the replica's override into the update.
3808    #[mz_ore::test]
3809    fn update_configuration_merges_replica_override() {
3810        let dyncfg = ConfigSet::default().add(&ENABLE_COLUMN_PAGED_BATCHER);
3811
3812        let replica = ReplicaId::User(1);
3813        let mut override_updates = ConfigUpdates::default();
3814        override_updates.add(&ENABLE_COLUMN_PAGED_BATCHER, true);
3815        let overrides = BTreeMap::from([(replica, override_updates)]);
3816
3817        let cmd = Instance::specialize_command_for_replica(
3818            ComputeCommand::UpdateConfiguration(Box::new(Default::default())),
3819            replica,
3820            &overrides,
3821            &dyncfg,
3822        );
3823        match cmd {
3824            ComputeCommand::UpdateConfiguration(params) => assert_eq!(
3825                params
3826                    .dyncfg_updates
3827                    .updates
3828                    .get(ENABLE_COLUMN_PAGED_BATCHER.name()),
3829                Some(&ConfigVal::Bool(true)),
3830            ),
3831            other => panic!("expected UpdateConfiguration, got {other:?}"),
3832        }
3833    }
3834}