Skip to main content

mz_compute_client/
controller.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! A controller that provides an interface to the compute layer, and the storage layer below it.
11//!
12//! The compute controller manages the creation, maintenance, and removal of compute instances.
13//! This involves ensuring the intended service state with the orchestrator, as well as maintaining
14//! a dedicated compute instance controller for each active compute instance.
15//!
16//! For each compute instance, the compute controller curates the creation of indexes and sinks
17//! installed on the instance, the progress of readers through these collections, and their
18//! eventual dropping and resource reclamation.
19//!
20//! The state maintained for a compute instance can be viewed as a partial map from `GlobalId` to
21//! collection. It is an error to use an identifier before it has been "created" with
22//! `create_dataflow()`. Once created, the controller holds a read capability for each output
23//! collection of a dataflow, which is manipulated with `set_read_policy()`. Eventually, a
24//! collection is dropped with `drop_collections()`.
25//!
26//! A dataflow can be in read-only or read-write mode. In read-only mode, the dataflow does not
27//! modify any persistent state. Sending a `allow_write` message to the compute instance will
28//! transition the dataflow to read-write mode, allowing it to write to persistent sinks.
29//!
30//! Created dataflows will prevent the compaction of their inputs, including other compute
31//! collections but also collections managed by the storage layer. Each dataflow input is prevented
32//! from compacting beyond the allowed compaction of each of its outputs, ensuring that we can
33//! recover each dataflow to its current state in case of failure or other reconfiguration.
34
35use std::collections::{BTreeMap, BTreeSet};
36use std::sync::{Arc, Mutex};
37use std::time::Duration;
38
39use mz_build_info::BuildInfo;
40use mz_cluster_client::client::ClusterReplicaLocation;
41use mz_cluster_client::metrics::ControllerMetrics;
42use mz_cluster_client::{ReplicaId, WallclockLagFn};
43use mz_compute_types::ComputeInstanceId;
44use mz_compute_types::config::ComputeReplicaConfig;
45use mz_compute_types::dataflows::DataflowDescription;
46use mz_compute_types::dyncfgs::{
47    COMPUTE_REPLICA_EXPIRATION_OFFSET, ENABLE_ARRANGEMENT_DICTIONARY_COMPRESSION_ALPHA,
48};
49use mz_dyncfg::{ConfigSet, ConfigUpdates};
50use mz_expr::RowSetFinishing;
51use mz_expr::row::RowCollection;
52use mz_ore::cast::CastFrom;
53use mz_ore::metrics::MetricsRegistry;
54use mz_ore::now::NowFn;
55use mz_ore::soft_assert_or_log;
56use mz_ore::soft_panic_or_log;
57use mz_ore::tracing::OpenTelemetryContext;
58use mz_persist_types::PersistLocation;
59use mz_repr::{GlobalId, RelationDesc, Row, Timestamp, frontier_within_lag};
60use mz_storage_client::controller::StorageController;
61use mz_storage_types::dyncfgs::ORE_OVERFLOWING_BEHAVIOR;
62use mz_storage_types::read_holds::ReadHold;
63use mz_storage_types::read_policy::ReadPolicy;
64use mz_storage_types::time_dependence::{TimeDependence, TimeDependenceError};
65use prometheus::proto::LabelPair;
66use serde::{Deserialize, Serialize};
67use timely::PartialOrder;
68use timely::progress::Antichain;
69use tokio::sync::{mpsc, oneshot};
70use tokio::time::{self, MissedTickBehavior};
71use uuid::Uuid;
72
73use crate::controller::error::{
74    CollectionLookupError, CollectionMissing, CollectionUpdateError, DataflowCreationError,
75    HydrationCheckBadTarget, InstanceExists, InstanceMissing, PeekError, ReadPolicyError,
76    ReplicaCreationError, ReplicaDropError,
77};
78use crate::controller::instance::{Instance, SharedCollectionState};
79use crate::controller::introspection::{IntrospectionUpdates, spawn_introspection_sink};
80use crate::controller::replica::ReplicaConfig;
81use crate::logging::{LogVariant, LoggingConfig};
82use crate::metrics::ComputeControllerMetrics;
83use crate::protocol::command::{ComputeParameters, PeekTarget};
84use crate::protocol::response::{PeekResponse, SubscribeBatch};
85
86mod instance;
87mod introspection;
88mod replica;
89mod sequential_hydration;
90
91pub mod error;
92pub mod instance_client;
93pub use instance_client::InstanceClient;
94
95pub(crate) type StorageCollections =
96    Arc<dyn mz_storage_client::storage_collections::StorageCollections + Send + Sync>;
97
98/// A collection's hydration and catch-up status.
99#[derive(Debug, Clone, Copy, PartialEq, Eq)]
100pub enum CollectionReadiness {
101    /// Hydrated and within the requested lag allowance.
102    Ready,
103    /// Hydrated, but behind the requested frontier.
104    Lagging {
105        /// The gap in timestamp ticks, or `None` when awaiting completion.
106        /// Ticks are milliseconds only on the epoch-milliseconds timeline.
107        lag: Option<u64>,
108    },
109    /// Not yet hydrated.
110    Unhydrated,
111}
112
113impl CollectionReadiness {
114    /// Classifies progress against an optional reference and lag allowance.
115    ///
116    /// The caller selects comparable frontiers: graceful replacement uses
117    /// per-replica output frontiers, whereas 0dt uses collection write frontiers.
118    /// No lag requirement means hydration alone. An empty reference requires
119    /// an empty frontier, since completion cannot be matched by finite progress.
120    pub fn classify(
121        hydrated: bool,
122        frontier: &Antichain<Timestamp>,
123        lag_requirement: Option<(&Antichain<Timestamp>, Timestamp)>,
124    ) -> Self {
125        if !hydrated {
126            Self::Unhydrated
127        } else if lag_requirement
128            .is_some_and(|(reference, lag)| !frontier_within_lag(frontier, reference, lag))
129        {
130            let reference = lag_requirement.expect("lag requirement was checked").0;
131            let lag =
132                frontier
133                    .as_option()
134                    .zip(reference.as_option())
135                    .map(|(frontier, reference)| {
136                        u64::from(*reference).saturating_sub(u64::from(*frontier))
137                    });
138            Self::Lagging { lag }
139        } else {
140            Self::Ready
141        }
142    }
143}
144
145/// Responses from the compute controller.
146#[derive(Debug)]
147pub enum ComputeControllerResponse {
148    /// See [`PeekNotification`].
149    PeekNotification(Uuid, PeekNotification, OpenTelemetryContext),
150    /// See [`crate::protocol::response::ComputeResponse::SubscribeResponse`].
151    SubscribeResponse(GlobalId, SubscribeBatch),
152    /// The response from a dataflow containing an `CopyToS3Oneshot` sink.
153    ///
154    /// The `GlobalId` identifies the sink. The `Result` is the response from
155    /// the sink, where an `Ok(n)` indicates that `n` rows were successfully
156    /// copied to S3 and an `Err` indicates that an error was encountered
157    /// during the copy operation.
158    ///
159    /// For a given `CopyToS3Oneshot` sink, there will be at most one `CopyToResponse`
160    /// produced. (The sink may produce no responses if its dataflow is dropped
161    /// before completion.)
162    CopyToResponse(GlobalId, Result<u64, anyhow::Error>),
163    /// A response reporting advancement of a collection's upper frontier.
164    ///
165    /// Once a collection's upper (aka "write frontier") has advanced to beyond a given time, the
166    /// contents of the collection as of that time have been sealed and cannot change anymore.
167    FrontierUpper {
168        /// The ID of a compute collection.
169        id: GlobalId,
170        /// The new upper frontier of the identified compute collection.
171        upper: Antichain<Timestamp>,
172    },
173}
174
175/// Notification and summary of a received and forwarded [`crate::protocol::response::ComputeResponse::PeekResponse`].
176#[derive(Clone, Debug, Serialize, Deserialize, PartialEq)]
177pub enum PeekNotification {
178    /// Returned rows of a successful peek.
179    Success {
180        /// Number of rows in the returned peek result.
181        rows: u64,
182        /// Size of the returned peek result in bytes.
183        result_size: u64,
184    },
185    /// Error of an unsuccessful peek, including the reason for the error.
186    Error(String),
187    /// The peek was canceled.
188    Canceled,
189}
190
191impl PeekNotification {
192    /// Construct a new [`PeekNotification`] from a [`PeekResponse`]. The `offset` and `limit`
193    /// parameters are used to calculate the number of rows in the peek result.
194    fn new(peek_response: &PeekResponse, offset: usize, limit: Option<usize>) -> Self {
195        match peek_response {
196            PeekResponse::Rows(rows) => {
197                let num_rows = u64::cast_from(RowCollection::offset_limit(
198                    rows.iter().map(|r| r.count()).sum(),
199                    offset,
200                    limit,
201                ));
202                let result_size = u64::cast_from(rows.iter().map(|r| r.byte_len()).sum::<usize>());
203
204                tracing::trace!(?num_rows, ?result_size, "inline peek result");
205
206                Self::Success {
207                    rows: num_rows,
208                    result_size,
209                }
210            }
211            PeekResponse::Stashed(stashed_response) => {
212                let rows = stashed_response.num_rows(offset, limit);
213                let result_size = stashed_response.size_bytes();
214
215                tracing::trace!(?rows, ?result_size, "stashed peek result");
216
217                Self::Success {
218                    rows: u64::cast_from(rows),
219                    result_size: u64::cast_from(result_size),
220                }
221            }
222            PeekResponse::Error(err) => Self::Error(err.to_string()),
223            PeekResponse::Canceled => Self::Canceled,
224        }
225    }
226}
227
228/// A controller for the compute layer.
229pub struct ComputeController {
230    instances: BTreeMap<ComputeInstanceId, InstanceState>,
231    /// A map from an instance ID to an arbitrary string that describes the
232    /// class of the workload that compute instance is running (e.g.,
233    /// `production` or `staging`).
234    instance_workload_classes: Arc<Mutex<BTreeMap<ComputeInstanceId, Option<String>>>>,
235    build_info: &'static BuildInfo,
236    /// A handle providing access to storage collections.
237    storage_collections: StorageCollections,
238    /// Set to `true` once `initialization_complete` has been called.
239    initialized: bool,
240    /// Whether or not this controller is in read-only mode.
241    ///
242    /// When in read-only mode, neither this controller nor the instances
243    /// controlled by it are allowed to affect changes to external systems
244    /// (largely persist).
245    read_only: bool,
246    /// Compute configuration to apply to new instances.
247    config: ComputeParameters,
248    /// The persist location where we can stash large peek results.
249    peek_stash_persist_location: PersistLocation,
250    /// A controller response to be returned on the next call to [`ComputeController::process`].
251    stashed_response: Option<ComputeControllerResponse>,
252    /// The compute controller metrics.
253    metrics: ComputeControllerMetrics,
254    /// A function that produces the current wallclock time.
255    now: NowFn,
256    /// A function that computes the lag between the given time and wallclock time.
257    wallclock_lag: WallclockLagFn<Timestamp>,
258    /// Dynamic system configuration.
259    ///
260    /// Updated through `ComputeController::update_configuration` calls and shared with all
261    /// subcomponents of the compute controller.
262    dyncfg: Arc<ConfigSet>,
263    /// The replica-local scoped overrides of [`Self::dyncfg`], by replica.
264    ///
265    /// Sparse, and kept here in addition to on the `Instance`s because replica
266    /// configuration that the controller resolves once, at replica creation,
267    /// must be read through the new replica's overrides.
268    replica_dyncfg_overrides: BTreeMap<ReplicaId, ConfigUpdates>,
269
270    /// Receiver for responses produced by `Instance`s.
271    response_rx: mpsc::UnboundedReceiver<ComputeControllerResponse>,
272    /// Response sender that's passed to new `Instance`s.
273    response_tx: mpsc::UnboundedSender<ComputeControllerResponse>,
274    /// Receiver for introspection updates produced by `Instance`s.
275    ///
276    /// When [`ComputeController::start_introspection_sink`] is first called, this receiver is
277    /// passed to the introspection sink task.
278    introspection_rx: Option<mpsc::UnboundedReceiver<IntrospectionUpdates>>,
279    /// Introspection updates sender that's passed to new `Instance`s.
280    introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
281
282    /// Ticker for scheduling periodic maintenance work.
283    maintenance_ticker: tokio::time::Interval,
284    /// Whether maintenance work was scheduled.
285    maintenance_scheduled: bool,
286}
287
288impl ComputeController {
289    /// Construct a new [`ComputeController`].
290    pub fn new(
291        build_info: &'static BuildInfo,
292        storage_collections: StorageCollections,
293        read_only: bool,
294        metrics_registry: &MetricsRegistry,
295        peek_stash_persist_location: PersistLocation,
296        controller_metrics: ControllerMetrics,
297        now: NowFn,
298        wallclock_lag: WallclockLagFn<Timestamp>,
299    ) -> Self {
300        let (response_tx, response_rx) = mpsc::unbounded_channel();
301        let (introspection_tx, introspection_rx) = mpsc::unbounded_channel();
302
303        let mut maintenance_ticker = time::interval(Duration::from_secs(1));
304        maintenance_ticker.set_missed_tick_behavior(MissedTickBehavior::Skip);
305
306        let instance_workload_classes = Arc::new(Mutex::new(BTreeMap::<
307            ComputeInstanceId,
308            Option<String>,
309        >::new()));
310
311        // Apply a `workload_class` label to all metrics in the registry that
312        // have an `instance_id` label for an instance whose workload class is
313        // known.
314        metrics_registry.register_postprocessor({
315            let instance_workload_classes = Arc::clone(&instance_workload_classes);
316            move |metrics| {
317                let instance_workload_classes = instance_workload_classes
318                    .lock()
319                    .expect("lock poisoned")
320                    .iter()
321                    .map(|(id, workload_class)| (id.to_string(), workload_class.clone()))
322                    .collect::<BTreeMap<String, Option<String>>>();
323                for metric in metrics {
324                    'metric: for metric in metric.mut_metric() {
325                        for label in metric.get_label() {
326                            if label.name() == "instance_id" {
327                                if let Some(workload_class) = instance_workload_classes
328                                    .get(label.value())
329                                    .cloned()
330                                    .flatten()
331                                {
332                                    let mut label = LabelPair::default();
333                                    label.set_name("workload_class".into());
334                                    label.set_value(workload_class.clone());
335
336                                    let mut labels = metric.take_label();
337                                    labels.push(label);
338                                    metric.set_label(labels);
339                                }
340                                continue 'metric;
341                            }
342                        }
343                    }
344                }
345            }
346        });
347
348        let metrics = ComputeControllerMetrics::new(metrics_registry, controller_metrics);
349
350        Self {
351            instances: BTreeMap::new(),
352            instance_workload_classes,
353            build_info,
354            storage_collections,
355            initialized: false,
356            read_only,
357            config: Default::default(),
358            peek_stash_persist_location,
359            stashed_response: None,
360            metrics,
361            now,
362            wallclock_lag,
363            dyncfg: Arc::new(mz_dyncfgs::all_dyncfgs()),
364            replica_dyncfg_overrides: BTreeMap::new(),
365            response_rx,
366            response_tx,
367            introspection_rx: Some(introspection_rx),
368            introspection_tx,
369            maintenance_ticker,
370            maintenance_scheduled: false,
371        }
372    }
373
374    /// Start sinking the compute controller's introspection data into storage.
375    ///
376    /// This method should be called once the introspection collections have been registered with
377    /// the storage controller. It will panic if invoked earlier than that.
378    pub fn start_introspection_sink(&mut self, storage_controller: &dyn StorageController) {
379        if let Some(rx) = self.introspection_rx.take() {
380            spawn_introspection_sink(rx, storage_controller);
381        }
382    }
383
384    /// TODO(database-issues#7533): Add documentation.
385    pub fn instance_exists(&self, id: ComputeInstanceId) -> bool {
386        self.instances.contains_key(&id)
387    }
388
389    /// Return a reference to the indicated compute instance.
390    fn instance(&self, id: ComputeInstanceId) -> Result<&InstanceState, InstanceMissing> {
391        self.instances.get(&id).ok_or(InstanceMissing(id))
392    }
393
394    /// Return an `InstanceClient` for the indicated compute instance.
395    pub fn instance_client(
396        &self,
397        id: ComputeInstanceId,
398    ) -> Result<InstanceClient, InstanceMissing> {
399        self.instance(id).map(|instance| instance.client.clone())
400    }
401
402    /// Return a mutable reference to the indicated compute instance.
403    fn instance_mut(
404        &mut self,
405        id: ComputeInstanceId,
406    ) -> Result<&mut InstanceState, InstanceMissing> {
407        self.instances.get_mut(&id).ok_or(InstanceMissing(id))
408    }
409
410    /// List the IDs of all collections in the identified compute instance.
411    pub fn collection_ids(
412        &self,
413        instance_id: ComputeInstanceId,
414    ) -> Result<impl Iterator<Item = GlobalId> + '_, InstanceMissing> {
415        let instance = self.instance(instance_id)?;
416        let ids = instance.collections.keys().copied();
417        Ok(ids)
418    }
419
420    /// Return the frontiers of the indicated collection.
421    ///
422    /// If an `instance_id` is provided, the collection is assumed to be installed on that
423    /// instance. Otherwise all available instances are searched.
424    pub fn collection_frontiers(
425        &self,
426        collection_id: GlobalId,
427        instance_id: Option<ComputeInstanceId>,
428    ) -> Result<CollectionFrontiers, CollectionLookupError> {
429        let collection = match instance_id {
430            Some(id) => self.instance(id)?.collection(collection_id)?,
431            None => self
432                .instances
433                .values()
434                .find_map(|i| i.collections.get(&collection_id))
435                .ok_or(CollectionMissing(collection_id))?,
436        };
437
438        Ok(collection.frontiers())
439    }
440
441    /// List compute collections that depend on the given collection.
442    pub fn collection_reverse_dependencies(
443        &self,
444        instance_id: ComputeInstanceId,
445        id: GlobalId,
446    ) -> Result<impl Iterator<Item = GlobalId> + '_, InstanceMissing> {
447        let instance = self.instance(instance_id)?;
448        let collections = instance.collections.iter();
449        let ids = collections
450            .filter_map(move |(cid, c)| c.compute_dependencies.contains(&id).then_some(*cid));
451        Ok(ids)
452    }
453
454    /// Returns `true` iff the given collection has been hydrated.
455    ///
456    /// For this check, zero-replica clusters are always considered hydrated.
457    /// Their collections would never normally be considered hydrated but it's
458    /// clearly intentional that they have no replicas.
459    pub async fn collection_hydrated(
460        &self,
461        instance_id: ComputeInstanceId,
462        collection_id: GlobalId,
463    ) -> Result<bool, anyhow::Error> {
464        let instance = self.instance(instance_id)?;
465
466        let res = instance
467            .call_sync(move |i| i.collection_hydrated(collection_id))
468            .await?;
469
470        Ok(res)
471    }
472
473    /// Returns `true` if all non-transient, non-excluded collections are ready on any of the
474    /// provided replicas: hydrated, and, when `allowed_lag` is `Some`, no further than that
475    /// behind the furthest output frontier any of the `reference_replicas` reports for the
476    /// collection.
477    ///
478    /// See `Instance::collections_ready_on_replicas` for why hydration alone is not a
479    /// readiness signal for a cut-over.
480    ///
481    /// For this check, zero-replica clusters are always considered ready.
482    /// Their collections would never normally be considered hydrated but it's
483    /// clearly intentional that they have no replicas.
484    pub fn collections_ready_for_replicas(
485        &self,
486        instance_id: ComputeInstanceId,
487        replicas: Vec<ReplicaId>,
488        exclude_collections: BTreeSet<GlobalId>,
489        allowed_lag: Option<Timestamp>,
490        reference_replicas: BTreeSet<ReplicaId>,
491    ) -> Result<oneshot::Receiver<bool>, anyhow::Error> {
492        let instance = self.instance(instance_id)?;
493
494        // Validation
495        if !instance.replicas.is_empty()
496            && !replicas.iter().any(|id| instance.replicas.contains(id))
497        {
498            return Err(HydrationCheckBadTarget(replicas).into());
499        }
500
501        let (tx, rx) = oneshot::channel();
502        instance.call(move |i| {
503            let result = i
504                .collections_ready_on_replicas(
505                    Some(replicas),
506                    &exclude_collections,
507                    allowed_lag,
508                    &reference_replicas,
509                )
510                .expect("validated");
511            let _ = tx.send(result);
512        });
513
514        Ok(rx)
515    }
516
517    /// Returns the state of the [`ComputeController`] formatted as JSON.
518    ///
519    /// The returned value is not guaranteed to be stable and may change at any point in time.
520    pub async fn dump(&self) -> Result<serde_json::Value, anyhow::Error> {
521        // Note: We purposefully use the `Debug` formatting for the value of all fields in the
522        // returned object as a tradeoff between usability and stability. `serde_json` will fail
523        // to serialize an object if the keys aren't strings, so `Debug` formatting the values
524        // prevents a future unrelated change from silently breaking this method.
525
526        // Destructure `self` here so we don't forget to consider dumping newly added fields.
527        let Self {
528            instances,
529            instance_workload_classes,
530            build_info: _,
531            storage_collections: _,
532            initialized,
533            read_only,
534            config: _,
535            peek_stash_persist_location: _,
536            stashed_response,
537            metrics: _,
538            now: _,
539            wallclock_lag: _,
540            dyncfg: _,
541            replica_dyncfg_overrides: _,
542            response_rx: _,
543            response_tx: _,
544            introspection_rx: _,
545            introspection_tx: _,
546            maintenance_ticker: _,
547            maintenance_scheduled,
548        } = self;
549
550        let mut instances_dump = BTreeMap::new();
551        for (id, instance) in instances {
552            let dump = instance.dump().await?;
553            instances_dump.insert(id.to_string(), dump);
554        }
555
556        let instance_workload_classes: BTreeMap<_, _> = instance_workload_classes
557            .lock()
558            .expect("lock poisoned")
559            .iter()
560            .map(|(id, wc)| (id.to_string(), format!("{wc:?}")))
561            .collect();
562
563        Ok(serde_json::json!({
564            "instances": instances_dump,
565            "instance_workload_classes": instance_workload_classes,
566            "initialized": initialized,
567            "read_only": read_only,
568            "stashed_response": format!("{stashed_response:?}"),
569            "maintenance_scheduled": maintenance_scheduled,
570        }))
571    }
572}
573
574impl ComputeController {
575    /// Create a compute instance.
576    pub fn create_instance(
577        &mut self,
578        id: ComputeInstanceId,
579        arranged_logs: BTreeMap<LogVariant, GlobalId>,
580        workload_class: Option<String>,
581    ) -> Result<(), InstanceExists> {
582        if self.instances.contains_key(&id) {
583            return Err(InstanceExists(id));
584        }
585
586        let mut collections = BTreeMap::new();
587        let mut logs = Vec::with_capacity(arranged_logs.len());
588        for (&log, &id) in &arranged_logs {
589            let collection = Collection::new_log();
590            let shared = collection.shared.clone();
591            collections.insert(id, collection);
592            logs.push((log, id, shared));
593        }
594
595        let client = InstanceClient::spawn(
596            id,
597            self.build_info,
598            Arc::clone(&self.storage_collections),
599            self.peek_stash_persist_location.clone(),
600            logs,
601            self.metrics.for_instance(id),
602            self.now.clone(),
603            self.wallclock_lag.clone(),
604            Arc::clone(&self.dyncfg),
605            self.response_tx.clone(),
606            self.introspection_tx.clone(),
607            self.read_only,
608        );
609
610        let instance = InstanceState::new(client, collections);
611        self.instances.insert(id, instance);
612
613        self.instance_workload_classes
614            .lock()
615            .expect("lock poisoned")
616            .insert(id, workload_class.clone());
617
618        let instance = self.instances.get_mut(&id).expect("instance just added");
619        if self.initialized {
620            instance.call(Instance::initialization_complete);
621        }
622
623        // The replica also receives the current dyncfg create-time, folded into `CreateInstance`
624        // so create-time setup observes synced values. This `UpdateConfiguration` is still
625        // required: it carries the rest of `ComputeParameters` (workload class, max result size,
626        // tracing) and syncs the dyncfg into the persist config and metrics, none of which ride in
627        // `CreateInstance`. The overlapping dyncfg application is idempotent.
628        let mut config_params = self.config.clone();
629        config_params.workload_class = Some(workload_class);
630        instance.call(|i| i.update_configuration(config_params));
631
632        Ok(())
633    }
634
635    /// Updates a compute instance's workload class.
636    pub fn update_instance_workload_class(
637        &mut self,
638        id: ComputeInstanceId,
639        workload_class: Option<String>,
640    ) -> Result<(), InstanceMissing> {
641        // Ensure that the instance exists first.
642        let _ = self.instance(id)?;
643
644        self.instance_workload_classes
645            .lock()
646            .expect("lock poisoned")
647            .insert(id, workload_class);
648
649        // Cause a config update to notify the instance about its new workload class.
650        self.update_configuration(Default::default());
651
652        Ok(())
653    }
654
655    /// Remove a compute instance.
656    ///
657    /// # Panics
658    ///
659    /// Panics if the identified `instance` still has active replicas.
660    pub fn drop_instance(&mut self, id: ComputeInstanceId) {
661        if let Some(instance) = self.instances.remove(&id) {
662            instance.call(|i| i.shutdown());
663        }
664
665        self.instance_workload_classes
666            .lock()
667            .expect("lock poisoned")
668            .remove(&id);
669    }
670
671    /// Returns the compute controller's config set.
672    pub fn dyncfg(&self) -> &Arc<ConfigSet> {
673        &self.dyncfg
674    }
675
676    /// Update compute configuration.
677    pub fn update_configuration(&mut self, config_params: ComputeParameters) {
678        // Apply dyncfg updates.
679        config_params.dyncfg_updates.apply(&self.dyncfg);
680
681        let instance_workload_classes = self
682            .instance_workload_classes
683            .lock()
684            .expect("lock poisoned");
685
686        // Forward updates to existing clusters.
687        // Workload classes are cluster-specific, so we need to overwrite them here.
688        for (id, instance) in self.instances.iter_mut() {
689            let mut params = config_params.clone();
690            params.workload_class = Some(instance_workload_classes[id].clone());
691            instance.call(|i| i.update_configuration(params));
692        }
693
694        let overflowing_behavior = ORE_OVERFLOWING_BEHAVIOR.get(&self.dyncfg);
695        match overflowing_behavior.parse() {
696            Ok(behavior) => mz_ore::overflowing::set_behavior(behavior),
697            Err(err) => {
698                tracing::error!(
699                    err,
700                    overflowing_behavior,
701                    "Invalid value for ore_overflowing_behavior"
702                );
703            }
704        }
705
706        // Remember updates for future clusters.
707        self.config.update(config_params);
708    }
709
710    /// Replaces the per-replica dyncfg overrides for the given instances.
711    ///
712    /// This only stores the overrides, here and on the instances; callers
713    /// should follow with a configuration push (e.g.
714    /// [`Self::update_configuration`]) so existing replicas observe the new
715    /// values. Instances absent from `overrides` have their overrides cleared,
716    /// so a replica that no longer has an override reverts to the
717    /// environment-wide configuration. Used by the scoped feature flags
718    /// (replica-local) layer.
719    pub fn update_replica_dyncfg_overrides(
720        &mut self,
721        mut overrides: BTreeMap<ComputeInstanceId, BTreeMap<ReplicaId, ConfigUpdates>>,
722    ) {
723        self.replica_dyncfg_overrides = overrides
724            .values()
725            .flat_map(|replicas| replicas.iter())
726            .map(|(replica_id, updates)| (*replica_id, updates.clone()))
727            .collect();
728        for (id, instance) in self.instances.iter_mut() {
729            let instance_overrides = overrides.remove(id).unwrap_or_default();
730            instance.call(move |i| i.update_replica_dyncfg_overrides(instance_overrides));
731        }
732    }
733
734    /// Mark the end of any initialization commands.
735    ///
736    /// The implementor may wait for this method to be called before implementing prior commands,
737    /// and so it is important for a user to invoke this method as soon as it is comfortable.
738    /// This method can be invoked immediately, at the potential expense of performance.
739    pub fn initialization_complete(&mut self) {
740        self.initialized = true;
741        for instance in self.instances.values_mut() {
742            instance.call(Instance::initialization_complete);
743        }
744    }
745
746    /// Wait until the controller is ready to do some processing.
747    ///
748    /// This method may block for an arbitrarily long time.
749    ///
750    /// When the method returns, the caller should call [`ComputeController::process`].
751    ///
752    /// This method is cancellation safe.
753    pub async fn ready(&mut self) {
754        if self.stashed_response.is_some() {
755            // We still have a response stashed, which we are immediately ready to process.
756            return;
757        }
758        if self.maintenance_scheduled {
759            // Maintenance work has been scheduled.
760            return;
761        }
762
763        tokio::select! {
764            resp = self.response_rx.recv() => {
765                let resp = resp.expect("`self.response_tx` not dropped");
766                self.stashed_response = Some(resp);
767            }
768            _ = self.maintenance_ticker.tick() => {
769                self.maintenance_scheduled = true;
770            },
771        }
772    }
773
774    /// Adds replicas of an instance.
775    pub fn add_replica_to_instance(
776        &mut self,
777        instance_id: ComputeInstanceId,
778        replica_id: ReplicaId,
779        location: ClusterReplicaLocation,
780        config: ComputeReplicaConfig,
781    ) -> Result<(), ReplicaCreationError> {
782        use ReplicaCreationError::*;
783
784        let instance = self.instance(instance_id)?;
785
786        // Validation
787        if instance.replicas.contains(&replica_id) {
788            return Err(ReplicaExists(replica_id));
789        }
790
791        let (enable_logging, interval) = match config.logging.interval {
792            Some(interval) => (true, interval),
793            None => (false, Duration::from_secs(1)),
794        };
795
796        // Both configs below are `ParameterScope::Replica` and are resolved
797        // here, once, for the replica being created. Reading them through the
798        // new replica's scoped overrides is what makes those declarations
799        // effective: the values are frozen into `ReplicaConfig` and never
800        // re-read from the environment-wide set. The overrides for a replica
801        // created by DDL are committed in the same transaction that creates it,
802        // so they are already installed by the time we get here.
803        let overrides = self.replica_dyncfg_overrides.get(&replica_id);
804
805        let expiration_offset =
806            COMPUTE_REPLICA_EXPIRATION_OFFSET.get_with_overrides(&self.dyncfg, overrides);
807
808        // Capture dictionary compression once, at replica creation, and hold it fixed for the
809        // replica's lifetime (see `InstanceConfig::arrangement_dictionary_compression`). This is
810        // why a later flip of the flag only affects replicas created afterwards. The feature flag
811        // only gates the feature: a replica honors its per-cluster/replica configured value only
812        // while the flag is enabled, so turning the flag off disables compression on new or
813        // restarted replicas regardless of their configuration.
814        let arrangement_dictionary_compression = ENABLE_ARRANGEMENT_DICTIONARY_COMPRESSION_ALPHA
815            .get_with_overrides(&self.dyncfg, overrides)
816            && config.arrangement_compression;
817
818        let replica_config = ReplicaConfig {
819            location,
820            logging: LoggingConfig {
821                interval,
822                enable_logging,
823                log_logging: config.logging.log_logging,
824                index_logs: Default::default(),
825            },
826            grpc_client: self.config.grpc_client.clone(),
827            expiration_offset: (!expiration_offset.is_zero()).then_some(expiration_offset),
828            arrangement_dictionary_compression,
829        };
830
831        let instance = self.instance_mut(instance_id).expect("validated");
832        instance.replicas.insert(replica_id);
833
834        instance.call(move |i| {
835            i.add_replica(replica_id, replica_config, None)
836                .expect("validated")
837        });
838
839        Ok(())
840    }
841
842    /// Removes a replica from an instance, including its service in the orchestrator.
843    pub fn drop_replica(
844        &mut self,
845        instance_id: ComputeInstanceId,
846        replica_id: ReplicaId,
847    ) -> Result<(), ReplicaDropError> {
848        use ReplicaDropError::*;
849
850        let instance = self.instance_mut(instance_id)?;
851
852        // Validation
853        if !instance.replicas.contains(&replica_id) {
854            return Err(ReplicaMissing(replica_id));
855        }
856
857        instance.replicas.remove(&replica_id);
858
859        // The coordinator only re-pushes the override map when the scoped
860        // configuration itself changes, so a dropped replica's entry would
861        // otherwise be retained until the next such change.
862        self.replica_dyncfg_overrides.remove(&replica_id);
863
864        let instance = self.instance_mut(instance_id).expect("validated");
865        instance.call(move |i| i.remove_replica(replica_id).expect("validated"));
866
867        Ok(())
868    }
869
870    /// Creates the described dataflow and initializes state for its output.
871    ///
872    /// Only sink exports are allowed to have a `target_replica`: materialized views, subscribes,
873    /// and metric sinks. A user's `CREATE METRIC SINK` runs untargeted, so every replica renders it
874    /// into its own registry. The coordinator's curated metric sinks are installed per replica and
875    /// do target one, so each replica's series are attributable to it.
876    ///
877    /// Panics if called with a dataflow description that has index exports
878    /// when `target_replica` is set.
879    pub fn create_dataflow(
880        &mut self,
881        instance_id: ComputeInstanceId,
882        mut dataflow: DataflowDescription<mz_compute_types::plan::LirRelationExpr, ()>,
883        target_replica: Option<ReplicaId>,
884    ) -> Result<(), DataflowCreationError> {
885        use DataflowCreationError::*;
886
887        let instance = self.instance(instance_id)?;
888
889        // Validation: target replica
890        if let Some(replica_id) = target_replica {
891            if !instance.replicas.contains(&replica_id) {
892                return Err(ReplicaMissing(replica_id));
893            }
894            assert!(
895                dataflow.exported_index_ids().next().is_none(),
896                "Replica-targeted indexes are not supported"
897            );
898        }
899
900        // Validation: as_of
901        let as_of = dataflow.as_of.as_ref().ok_or(MissingAsOf)?;
902        if as_of.is_empty() && dataflow.subscribe_ids().next().is_some() {
903            return Err(EmptyAsOfForSubscribe);
904        }
905        if as_of.is_empty() && dataflow.copy_to_ids().next().is_some() {
906            return Err(EmptyAsOfForCopyTo);
907        }
908
909        // Validation: the dataflow exports something
910        //
911        // An export-less description has nothing to render and no answer to "what do the exports
912        // read", which the checks below are phrased in terms of. `optimize_dataflow` leaves such a
913        // description's imports alone for that reason, so one arriving here would fail the import
914        // check for the wrong reason.
915        soft_assert_or_log!(
916            !dataflow.index_exports.is_empty() || !dataflow.sink_exports.is_empty(),
917            "dataflow {} has no exports",
918            dataflow.debug_name,
919        );
920
921        // The imports the exports actually read. `optimize_dataflow` prunes the import list to
922        // exactly this set, so the two agree unless a producer stopped pruning.
923        //
924        // Computed once and used twice: the check below reports a loose list, and
925        // `determine_time_dependence` counts through it rather than over the raw list. That
926        // consumer is the one whose wrong answer hangs an environment: an import no export reads
927        // would report wall-clock dependence for a dataflow whose exports are constant, earning it
928        // a dataflow expiration that pins the output frontier days short of the empty antichain,
929        // and nothing downstream would learn the collection is final. Deriving it from this set
930        // makes that correct by construction, leaving the prune to reclaim the read hold and the
931        // persist source.
932        let used_imports = dataflow.used_import_ids();
933
934        // Validation: every import is read
935        //
936        // The read holds and the persist sources the replicas build are still derived from the raw
937        // list below, so a loose one describes a dataflow other than the one that will run. A
938        // logging variant rather than `soft_assert_no_log!`: the walk is paid for above either way,
939        // so reporting it in production costs only the comparison.
940        soft_assert_or_log!(
941            dataflow.import_ids().all(|id| used_imports.contains(&id)),
942            "dataflow {} imports collections no export reads: imports {:?}, read {:?}",
943            dataflow.debug_name,
944            dataflow.import_ids().collect::<Vec<_>>(),
945            used_imports,
946        );
947
948        // Validation: input collections
949        let storage_ids = dataflow.imported_source_ids().collect();
950        let mut import_read_holds = self.storage_collections.acquire_read_holds(storage_ids)?;
951        for id in dataflow.imported_index_ids() {
952            let read_hold = instance.acquire_read_hold(id)?;
953            import_read_holds.push(read_hold);
954        }
955        for hold in &import_read_holds {
956            if PartialOrder::less_than(as_of, hold.since()) {
957                return Err(SinceViolation(hold.id()));
958            }
959        }
960
961        // Validation: storage sink collections
962        for id in dataflow.persist_sink_ids() {
963            if self.storage_collections.check_exists(id).is_err() {
964                return Err(CollectionMissing(id));
965            }
966        }
967        let time_dependence = self
968            .determine_time_dependence(instance_id, &dataflow, &used_imports)
969            .expect("must exist");
970
971        let instance = self.instance_mut(instance_id).expect("validated");
972
973        let mut shared_collection_state = BTreeMap::new();
974        for id in dataflow.export_ids() {
975            let shared = SharedCollectionState::new(as_of.clone());
976            let collection = Collection {
977                write_only: dataflow.sink_exports.contains_key(&id),
978                compute_dependencies: dataflow.imported_index_ids().collect(),
979                shared: shared.clone(),
980                time_dependence: time_dependence.clone(),
981            };
982            instance.collections.insert(id, collection);
983            shared_collection_state.insert(id, shared);
984        }
985
986        dataflow.time_dependence = time_dependence;
987
988        instance.call(move |i| {
989            i.create_dataflow(
990                dataflow,
991                import_read_holds,
992                shared_collection_state,
993                target_replica,
994            )
995            .expect("validated")
996        });
997
998        Ok(())
999    }
1000
1001    /// Drop the read capability for the given collections and allow their resources to be
1002    /// reclaimed.
1003    pub fn drop_collections(
1004        &mut self,
1005        instance_id: ComputeInstanceId,
1006        collection_ids: Vec<GlobalId>,
1007    ) -> Result<(), CollectionUpdateError> {
1008        let instance = self.instance_mut(instance_id)?;
1009
1010        // Validation
1011        for id in &collection_ids {
1012            instance.collection(*id)?;
1013        }
1014
1015        for id in &collection_ids {
1016            instance.collections.remove(id);
1017        }
1018
1019        instance.call(|i| i.drop_collections(collection_ids).expect("validated"));
1020
1021        Ok(())
1022    }
1023
1024    /// Initiate a peek request for the contents of the given collection at `timestamp`.
1025    ///
1026    /// The caller supplies a `read_hold` for the peek target — via
1027    /// [`ComputeController::acquire_read_hold`] for `PeekTarget::Index`, or via the storage
1028    /// collections for `PeekTarget::Persist`. The hold keeps the collection's `since` at
1029    /// `<= timestamp` until the peek completes.
1030    pub fn peek(
1031        &self,
1032        instance_id: ComputeInstanceId,
1033        peek_target: PeekTarget,
1034        literal_constraints: Option<Vec<Row>>,
1035        uuid: Uuid,
1036        timestamp: Timestamp,
1037        result_desc: RelationDesc,
1038        finishing: RowSetFinishing,
1039        map_filter_project: mz_expr::SafeMfpPlan,
1040        read_hold: ReadHold,
1041        target_replica: Option<ReplicaId>,
1042        peek_response_tx: oneshot::Sender<PeekResponse>,
1043    ) -> Result<(), PeekError> {
1044        use PeekError::*;
1045
1046        let instance = self.instance(instance_id)?;
1047
1048        // Validation: target replica
1049        if let Some(replica_id) = target_replica {
1050            if !instance.replicas.contains(&replica_id) {
1051                return Err(ReplicaMissing(replica_id));
1052            }
1053        }
1054
1055        // Validation: the read hold must target this collection and must hold its `since`
1056        // at `<= timestamp`.
1057        if read_hold.id() != peek_target.id() {
1058            return Err(ReadHoldIdMismatch(read_hold.id()));
1059        }
1060        if !read_hold.since().less_equal(&timestamp) {
1061            return Err(SinceViolation(peek_target.id()));
1062        }
1063
1064        instance.call(move |i| {
1065            i.peek(
1066                peek_target,
1067                literal_constraints,
1068                uuid,
1069                timestamp,
1070                result_desc,
1071                finishing,
1072                map_filter_project,
1073                read_hold,
1074                target_replica,
1075                peek_response_tx,
1076            )
1077            .expect("validated")
1078        });
1079
1080        Ok(())
1081    }
1082
1083    /// Cancel an existing peek request.
1084    ///
1085    /// Canceling a peek is best effort. The caller may see any of the following
1086    /// after canceling a peek request:
1087    ///
1088    ///   * A `PeekResponse::Rows` indicating that the cancellation request did
1089    ///     not take effect in time and the query succeeded.
1090    ///   * A `PeekResponse::Canceled` affirming that the peek was canceled.
1091    ///   * No `PeekResponse` at all.
1092    pub fn cancel_peek(
1093        &self,
1094        instance_id: ComputeInstanceId,
1095        uuid: Uuid,
1096        reason: PeekResponse,
1097    ) -> Result<(), InstanceMissing> {
1098        self.instance(instance_id)?
1099            .call(move |i| i.cancel_peek(uuid, reason));
1100        Ok(())
1101    }
1102
1103    /// Assign a read policy to specific identifiers.
1104    ///
1105    /// The policies are assigned in the order presented, and repeated identifiers should
1106    /// conclude with the last policy. Changing a policy will immediately downgrade the read
1107    /// capability if appropriate, but it will not "recover" the read capability if the prior
1108    /// capability is already ahead of it.
1109    ///
1110    /// Identifiers not present in `policies` retain their existing read policies.
1111    ///
1112    /// It is an error to attempt to set a read policy for a collection that is not readable in the
1113    /// context of compute. At this time, only indexes are readable compute collections.
1114    pub fn set_read_policy(
1115        &self,
1116        instance_id: ComputeInstanceId,
1117        policies: Vec<(GlobalId, ReadPolicy)>,
1118    ) -> Result<(), ReadPolicyError> {
1119        use ReadPolicyError::*;
1120
1121        let instance = self.instance(instance_id)?;
1122
1123        // Validation
1124        for (id, _) in &policies {
1125            let collection = instance.collection(*id)?;
1126            if collection.write_only {
1127                return Err(WriteOnlyCollection(*id));
1128            }
1129        }
1130
1131        self.instance(instance_id)?
1132            .call(|i| i.set_read_policy(policies).expect("validated"));
1133
1134        Ok(())
1135    }
1136
1137    /// Acquires a [`ReadHold`] for the identified compute collection.
1138    pub fn acquire_read_hold(
1139        &self,
1140        instance_id: ComputeInstanceId,
1141        collection_id: GlobalId,
1142    ) -> Result<ReadHold, CollectionUpdateError> {
1143        let read_hold = self
1144            .instance(instance_id)?
1145            .acquire_read_hold(collection_id)?;
1146        Ok(read_hold)
1147    }
1148
1149    /// Determine the time dependence for a dataflow.
1150    ///
1151    /// `used_imports` are the imports the exports read, as
1152    /// [`DataflowDescription::used_import_ids`] reports them. Only those count: an import no export
1153    /// reads would report wall-clock dependence for a dataflow whose exports are constant, and that
1154    /// earns it a dataflow expiration, which pins its output frontier at the expiration time. A
1155    /// constant export's frontier is the empty antichain, so the pin would hold it days short of
1156    /// the truth and whoever reads that frontier would never learn the collection can no longer
1157    /// change.
1158    ///
1159    /// `optimize_dataflow` prunes the import list to this set, so the two agree and the filtering
1160    /// is a no-op. It is here because this is the consumer whose wrong answer hangs an environment,
1161    /// and deriving the answer from the read set makes it independent of the list staying tight.
1162    fn determine_time_dependence(
1163        &self,
1164        instance_id: ComputeInstanceId,
1165        dataflow: &DataflowDescription<mz_compute_types::plan::LirRelationExpr, ()>,
1166        used_imports: &BTreeSet<GlobalId>,
1167    ) -> Result<Option<TimeDependence>, TimeDependenceError> {
1168        let instance = self
1169            .instance(instance_id)
1170            .map_err(|err| TimeDependenceError::InstanceMissing(err.0))?;
1171        let mut time_dependencies = Vec::new();
1172
1173        for id in dataflow
1174            .imported_index_ids()
1175            .filter(|id| used_imports.contains(id))
1176        {
1177            let dependence = instance
1178                .get_time_dependence(id)
1179                .map_err(|err| TimeDependenceError::CollectionMissing(err.0))?;
1180            time_dependencies.push(dependence);
1181        }
1182
1183        'source: for id in dataflow
1184            .imported_source_ids()
1185            .filter(|id| used_imports.contains(id))
1186        {
1187            // We first check whether the id is backed by a compute object, in which case we use
1188            // the time dependence we know. This is true for storage sinks.
1189            for instance in self.instances.values() {
1190                if let Ok(dependence) = instance.get_time_dependence(id) {
1191                    time_dependencies.push(dependence);
1192                    continue 'source;
1193                }
1194            }
1195
1196            // Not a compute object: Consult the storage collections controller.
1197            time_dependencies.push(self.storage_collections.determine_time_dependence(id)?);
1198        }
1199
1200        Ok(TimeDependence::merge(
1201            time_dependencies,
1202            dataflow.refresh_schedule.as_ref(),
1203        ))
1204    }
1205
1206    /// Processes the work queued by [`ComputeController::ready`].
1207    #[mz_ore::instrument(level = "debug")]
1208    pub fn process(&mut self) -> Option<ComputeControllerResponse> {
1209        // Perform periodic maintenance work.
1210        if self.maintenance_scheduled {
1211            self.maintain();
1212            self.maintenance_scheduled = false;
1213        }
1214
1215        // Return a ready response, if any.
1216        self.stashed_response.take()
1217    }
1218
1219    #[mz_ore::instrument(level = "debug")]
1220    fn maintain(&mut self) {
1221        // Perform instance maintenance work.
1222        for instance in self.instances.values_mut() {
1223            instance.call(Instance::maintain);
1224        }
1225    }
1226
1227    /// Allow writes for the specified collections on `instance_id`.
1228    ///
1229    /// If this controller is in read-only mode, this is a no-op.
1230    pub fn allow_writes(
1231        &mut self,
1232        instance_id: ComputeInstanceId,
1233        collection_id: GlobalId,
1234    ) -> Result<(), CollectionUpdateError> {
1235        if self.read_only {
1236            tracing::debug!("Skipping allow_writes in read-only mode");
1237            return Ok(());
1238        }
1239
1240        self.allow_writes_inner(instance_id, collection_id)
1241    }
1242
1243    /// Like [`Self::allow_writes`], but takes effect even in read-only mode.
1244    ///
1245    /// The caller must guarantee that no leader environment writes the collection's output shard.
1246    /// In a 0dt deployment that means a shard this environment created for itself, the replacement
1247    /// shard of a `Replacement`-migrated builtin collection, never one the leader is still serving
1248    /// from. `Evolution` migrates in place and reuses the leader's shard, so it must not reach this
1249    /// path. Violating the guarantee races two writers on one shard.
1250    ///
1251    /// NOTE: ownership is exclusive per (build version, deploy generation), not per process: the
1252    /// migration shard entry naming the shard is keyed by that pair and a read-only catalog open is
1253    /// a savepoint, so two read-only processes of one generation both write it. Same shape as a
1254    /// multi-replica materialized view, which the self-correcting persist sink tolerates (see the
1255    /// `mz_compute::sink::materialized_view` module docs).
1256    ///
1257    /// This is the compute-side counterpart to the storage controller's `force_writable` handling
1258    /// of migrated storage collections. Migrated builtin tables are storage collections that
1259    /// storage force-writes read-only; migrated builtin MVs are compute collections that only this
1260    /// path can force-write. Both rest on the same guarantee (this environment exclusively owns the
1261    /// replacement shard) but run on separate write paths, so each needs its own bypass.
1262    ///
1263    /// NOTE: the replica-side handler enables persist compaction process-wide on the clusterd
1264    /// (`ComputeState::handle_allow_writes`), which this path is the first to trigger inside a
1265    /// read-only deployment.
1266    pub fn allow_writes_in_read_only(
1267        &mut self,
1268        instance_id: ComputeInstanceId,
1269        collection_id: GlobalId,
1270    ) -> Result<(), CollectionUpdateError> {
1271        // Every builtin eligible for this bypass has a system id, so a non-system id means the
1272        // caller's `Replacement`-only invariant broke. No-op rather than risk writing a shard the
1273        // leader still serves. The collection then sits in the caught-up gate on an unwritten
1274        // shard and blocks promotion, which is the loud, safe direction to fail.
1275        //
1276        // Storage asserts the same invariant hard, in
1277        // `StorageController::register_introspection_collection`. The asymmetry is deliberate: a
1278        // soft panic keeps a caller bug visible in CI and Sentry without downing production.
1279        if self.read_only && !collection_id.is_system() {
1280            soft_panic_or_log!(
1281                "allow_writes_in_read_only called for non-system collection {collection_id}; \
1282                 falling back to read-only no-op"
1283            );
1284            return Ok(());
1285        }
1286
1287        self.allow_writes_inner(instance_id, collection_id)
1288    }
1289
1290    fn allow_writes_inner(
1291        &mut self,
1292        instance_id: ComputeInstanceId,
1293        collection_id: GlobalId,
1294    ) -> Result<(), CollectionUpdateError> {
1295        let instance = self.instance_mut(instance_id)?;
1296
1297        // Validation
1298        instance.collection(collection_id)?;
1299
1300        instance.call(move |i| i.allow_writes(collection_id).expect("validated"));
1301
1302        Ok(())
1303    }
1304}
1305
1306#[derive(Debug)]
1307struct InstanceState {
1308    client: InstanceClient,
1309    replicas: BTreeSet<ReplicaId>,
1310    collections: BTreeMap<GlobalId, Collection>,
1311}
1312
1313impl InstanceState {
1314    fn new(client: InstanceClient, collections: BTreeMap<GlobalId, Collection>) -> Self {
1315        Self {
1316            client,
1317            replicas: Default::default(),
1318            collections,
1319        }
1320    }
1321
1322    fn collection(&self, id: GlobalId) -> Result<&Collection, CollectionMissing> {
1323        self.collections.get(&id).ok_or(CollectionMissing(id))
1324    }
1325
1326    /// Calls the given function on the instance task. Does not await the result.
1327    ///
1328    /// # Panics
1329    ///
1330    /// Panics if the instance corresponding to `self` does not exist.
1331    fn call<F>(&self, f: F)
1332    where
1333        F: FnOnce(&mut Instance) + Send + 'static,
1334    {
1335        self.client.call(f).expect("instance not dropped")
1336    }
1337
1338    /// Calls the given function on the instance task, and awaits the result.
1339    ///
1340    /// # Panics
1341    ///
1342    /// Panics if the instance corresponding to `self` does not exist.
1343    async fn call_sync<F, R>(&self, f: F) -> R
1344    where
1345        F: FnOnce(&mut Instance) -> R + Send + 'static,
1346        R: Send + 'static,
1347    {
1348        self.client
1349            .call_sync(f)
1350            .await
1351            .expect("instance not dropped")
1352    }
1353
1354    /// Acquires a [`ReadHold`] for the identified compute collection.
1355    pub fn acquire_read_hold(&self, id: GlobalId) -> Result<ReadHold, CollectionMissing> {
1356        // We acquire read holds at the earliest possible time rather than returning a copy
1357        // of the implied read hold. This is so that in `create_dataflow` we can acquire read holds
1358        // on compute dependencies at frontiers that are held back by other read holds the caller
1359        // has previously taken.
1360        //
1361        // If/when we change the compute API to expect callers to pass in the `ReadHold`s rather
1362        // than acquiring them ourselves, we might tighten this up and instead acquire read holds
1363        // at the implied capability.
1364
1365        let collection = self.collection(id)?;
1366        let since = collection.shared.lock_read_capabilities(|caps| {
1367            let since = caps.frontier().to_owned();
1368            caps.update_iter(since.iter().map(|t| (t.clone(), 1)));
1369            since
1370        });
1371
1372        let hold = ReadHold::new(id, since, self.client.read_hold_tx());
1373        Ok(hold)
1374    }
1375
1376    /// Return the stored time dependence for a collection.
1377    fn get_time_dependence(
1378        &self,
1379        id: GlobalId,
1380    ) -> Result<Option<TimeDependence>, CollectionMissing> {
1381        Ok(self.collection(id)?.time_dependence.clone())
1382    }
1383
1384    /// Returns the [`InstanceState`] formatted as JSON.
1385    pub async fn dump(&self) -> Result<serde_json::Value, anyhow::Error> {
1386        // Destructure `self` here so we don't forget to consider dumping newly added fields.
1387        let Self {
1388            client: _,
1389            replicas,
1390            collections,
1391        } = self;
1392
1393        let instance = self.call_sync(|i| i.dump()).await?;
1394        let replicas: Vec<_> = replicas.iter().map(|id| id.to_string()).collect();
1395        let collections: BTreeMap<_, _> = collections
1396            .iter()
1397            .map(|(id, c)| (id.to_string(), format!("{c:?}")))
1398            .collect();
1399
1400        Ok(serde_json::json!({
1401            "instance": instance,
1402            "replicas": replicas,
1403            "collections": collections,
1404        }))
1405    }
1406}
1407
1408#[derive(Debug)]
1409struct Collection {
1410    /// Whether a collection is write-only, i.e., we cannot read it directly like an index.
1411    write_only: bool,
1412    compute_dependencies: BTreeSet<GlobalId>,
1413    shared: SharedCollectionState,
1414    /// The computed time dependence for this collection. None indicates no specific information,
1415    /// a value describes how the collection relates to wall-clock time.
1416    time_dependence: Option<TimeDependence>,
1417}
1418
1419impl Collection {
1420    fn new_log() -> Self {
1421        let as_of = Antichain::from_elem(Timestamp::MIN);
1422        Self {
1423            write_only: false,
1424            compute_dependencies: Default::default(),
1425            shared: SharedCollectionState::new(as_of),
1426            time_dependence: Some(TimeDependence::default()),
1427        }
1428    }
1429
1430    fn frontiers(&self) -> CollectionFrontiers {
1431        let read_frontier = self
1432            .shared
1433            .lock_read_capabilities(|c| c.frontier().to_owned());
1434        let write_frontier = self.shared.lock_write_frontier(|f| f.clone());
1435        CollectionFrontiers {
1436            read_frontier,
1437            write_frontier,
1438        }
1439    }
1440}
1441
1442/// The frontiers of a compute collection.
1443#[derive(Clone, Debug)]
1444pub struct CollectionFrontiers {
1445    /// The read frontier.
1446    pub read_frontier: Antichain<Timestamp>,
1447    /// The write frontier.
1448    pub write_frontier: Antichain<Timestamp>,
1449}
1450
1451impl Default for CollectionFrontiers {
1452    fn default() -> Self {
1453        Self {
1454            read_frontier: Antichain::from_elem(Timestamp::MIN),
1455            write_frontier: Antichain::from_elem(Timestamp::MIN),
1456        }
1457    }
1458}