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(×tamp) {
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}