1use std::collections::{BTreeMap, BTreeSet};
13use std::fmt::Debug;
14use std::sync::{Arc, Mutex};
15use std::time::{Duration, Instant};
16
17use chrono::{DateTime, DurationRound, TimeDelta, Utc};
18use differential_dataflow::lattice::Lattice;
19use mz_build_info::BuildInfo;
20use mz_cluster_client::WallclockLagFn;
21use mz_compute_types::dataflows::{BuildDesc, DataflowDescription};
22use mz_compute_types::plan::render_plan::RenderPlan;
23use mz_compute_types::sinks::{
24 ComputeSinkConnection, ComputeSinkDesc, MaterializedViewSinkConnection,
25};
26use mz_compute_types::sources::SourceInstanceDesc;
27use mz_controller_types::dyncfgs::{
28 ENABLE_PAUSED_CLUSTER_READHOLD_DOWNGRADE, WALLCLOCK_LAG_RECORDING_INTERVAL,
29};
30use mz_dyncfg::{ConfigSet, ConfigUpdates};
31use mz_expr::RowSetFinishing;
32use mz_ore::cast::CastFrom;
33use mz_ore::channel::instrumented_unbounded_channel;
34use mz_ore::now::NowFn;
35use mz_ore::tracing::OpenTelemetryContext;
36use mz_ore::{soft_assert_or_log, soft_panic_or_log};
37use mz_persist_types::PersistLocation;
38use mz_repr::adt::timestamp::CheckedTimestamp;
39use mz_repr::refresh_schedule::RefreshSchedule;
40use mz_repr::{Datum, Diff, GlobalId, RelationDesc, Row, Timestamp};
41use mz_storage_client::controller::{IntrospectionType, WallclockLag, WallclockLagHistogramPeriod};
42use mz_storage_types::read_holds::{self, ReadHold};
43use mz_storage_types::read_policy::ReadPolicy;
44use thiserror::Error;
45use timely::PartialOrder;
46use timely::progress::frontier::MutableAntichain;
47use timely::progress::{Antichain, ChangeBatch};
48use tokio::sync::{mpsc, oneshot};
49use uuid::Uuid;
50
51use crate::controller::error::{
52 CollectionMissing, ERROR_TARGET_REPLICA_FAILED, HydrationCheckBadTarget,
53};
54use crate::controller::instance_client::PeekError;
55use crate::controller::replica::{ReplicaClient, ReplicaConfig};
56use crate::controller::{
57 CollectionReadiness, ComputeControllerResponse, IntrospectionUpdates, PeekNotification,
58 ReplicaId, StorageCollections,
59};
60use crate::logging::LogVariant;
61use crate::metrics::IntCounter;
62use crate::metrics::{InstanceMetrics, ReplicaCollectionMetrics, ReplicaMetrics, UIntGauge};
63use crate::protocol::command::{
64 ComputeCommand, ComputeParameters, InstanceConfig, Peek, PeekTarget,
65};
66use crate::protocol::history::ComputeCommandHistory;
67use crate::protocol::response::{
68 ComputeResponse, CopyToResponse, FrontiersResponse, PeekError as ProtocolPeekError,
69 PeekResponse, StatusResponse, SubscribeBatch, SubscribeResponse,
70};
71
72#[derive(Error, Debug)]
73#[error("replica exists already: {0}")]
74pub(super) struct ReplicaExists(pub ReplicaId);
75
76#[derive(Error, Debug)]
77#[error("replica does not exist: {0}")]
78pub(super) struct ReplicaMissing(pub ReplicaId);
79
80#[derive(Error, Debug)]
81pub(super) enum DataflowCreationError {
82 #[error("collection does not exist: {0}")]
83 CollectionMissing(GlobalId),
84 #[error("replica does not exist: {0}")]
85 ReplicaMissing(ReplicaId),
86 #[error("dataflow definition lacks an as_of value")]
87 MissingAsOf,
88 #[error("subscribe dataflow has an empty as_of")]
89 EmptyAsOfForSubscribe,
90 #[error("copy to dataflow has an empty as_of")]
91 EmptyAsOfForCopyTo,
92 #[error("no read hold provided for dataflow import: {0}")]
93 ReadHoldMissing(GlobalId),
94 #[error("insufficient read hold provided for dataflow import: {0}")]
95 ReadHoldInsufficient(GlobalId),
96}
97
98impl From<CollectionMissing> for DataflowCreationError {
99 fn from(error: CollectionMissing) -> Self {
100 Self::CollectionMissing(error.0)
101 }
102}
103
104#[derive(Error, Debug)]
105pub(super) enum ReadPolicyError {
106 #[error("collection does not exist: {0}")]
107 CollectionMissing(GlobalId),
108 #[error("collection is write-only: {0}")]
109 WriteOnlyCollection(GlobalId),
110}
111
112impl From<CollectionMissing> for ReadPolicyError {
113 fn from(error: CollectionMissing) -> Self {
114 Self::CollectionMissing(error.0)
115 }
116}
117
118pub(super) type Command = Box<dyn FnOnce(&mut Instance) + Send>;
120
121pub(super) type ReplicaResponse = (ReplicaId, u64, ComputeResponse);
124
125pub(super) struct Instance {
127 build_info: &'static BuildInfo,
129 storage_collections: StorageCollections,
131 initialized: bool,
133 read_only: bool,
138 workload_class: Option<String>,
142 replicas: BTreeMap<ReplicaId, ReplicaState>,
144 replica_dyncfg_overrides: BTreeMap<ReplicaId, ConfigUpdates>,
151 collections: BTreeMap<GlobalId, CollectionState>,
159 log_sources: BTreeMap<LogVariant, GlobalId>,
161 peeks: BTreeMap<Uuid, PendingPeek>,
170 subscribes: BTreeMap<GlobalId, ActiveSubscribe>,
184 copy_tos: BTreeSet<GlobalId>,
192 history: ComputeCommandHistory<UIntGauge>,
194 command_rx: mpsc::UnboundedReceiver<Command>,
196 response_tx: mpsc::UnboundedSender<ComputeControllerResponse>,
198 introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
200 metrics: InstanceMetrics,
202 dyncfg: Arc<ConfigSet>,
204
205 peek_stash_persist_location: PersistLocation,
207
208 now: NowFn,
210 wallclock_lag: WallclockLagFn<Timestamp>,
212 wallclock_lag_last_recorded: DateTime<Utc>,
214
215 read_hold_tx: read_holds::ChangeTx,
220 replica_tx: mz_ore::channel::InstrumentedUnboundedSender<ReplicaResponse, IntCounter>,
222 replica_rx: mz_ore::channel::InstrumentedUnboundedReceiver<ReplicaResponse, IntCounter>,
224}
225
226impl Instance {
227 fn collection(&self, id: GlobalId) -> Result<&CollectionState, CollectionMissing> {
229 self.collections.get(&id).ok_or(CollectionMissing(id))
230 }
231
232 fn collection_mut(&mut self, id: GlobalId) -> Result<&mut CollectionState, CollectionMissing> {
234 self.collections.get_mut(&id).ok_or(CollectionMissing(id))
235 }
236
237 fn expect_collection(&self, id: GlobalId) -> &CollectionState {
243 self.collections.get(&id).expect("collection must exist")
244 }
245
246 fn expect_collection_mut(&mut self, id: GlobalId) -> &mut CollectionState {
252 self.collections
253 .get_mut(&id)
254 .expect("collection must exist")
255 }
256
257 fn collections_iter(&self) -> impl Iterator<Item = (GlobalId, &CollectionState)> {
258 self.collections.iter().map(|(id, coll)| (*id, coll))
259 }
260
261 fn replicas_hosting(
268 &self,
269 id: GlobalId,
270 ) -> Result<impl Iterator<Item = &ReplicaState> + Clone, CollectionMissing> {
271 let target = self.collection(id)?.target_replica;
272 Ok(self
273 .replicas
274 .values()
275 .filter(move |r| target.map_or(true, |t| t == r.id)))
276 }
277
278 fn add_collection(
284 &mut self,
285 id: GlobalId,
286 as_of: Antichain<Timestamp>,
287 shared: SharedCollectionState,
288 storage_dependencies: BTreeMap<GlobalId, ReadHold>,
289 compute_dependencies: BTreeMap<GlobalId, ReadHold>,
290 replica_input_read_holds: Vec<ReadHold>,
291 write_only: bool,
292 storage_sink: bool,
293 initial_as_of: Option<Antichain<Timestamp>>,
294 refresh_schedule: Option<RefreshSchedule>,
295 target_replica: Option<ReplicaId>,
296 ) {
297 let dependency_ids: Vec<GlobalId> = compute_dependencies
299 .keys()
300 .chain(storage_dependencies.keys())
301 .copied()
302 .collect();
303 let introspection = CollectionIntrospection::new(
304 id,
305 self.introspection_tx.clone(),
306 as_of.clone(),
307 storage_sink,
308 initial_as_of,
309 refresh_schedule,
310 dependency_ids,
311 );
312 let mut state = CollectionState::new(
313 id,
314 as_of.clone(),
315 shared,
316 storage_dependencies,
317 compute_dependencies,
318 Arc::clone(&self.read_hold_tx),
319 introspection,
320 );
321 state.target_replica = target_replica;
322 if write_only {
324 state.read_policy = None;
325 }
326
327 if let Some(previous) = self.collections.insert(id, state) {
328 panic!("attempt to add a collection with existing ID {id} (previous={previous:?}");
329 }
330
331 for replica in self.replicas.values_mut() {
333 if target_replica.is_some_and(|id| id != replica.id) {
334 continue;
335 }
336 replica.add_collection(id, as_of.clone(), replica_input_read_holds.clone());
337 }
338 }
339
340 fn remove_collection(&mut self, id: GlobalId) {
341 for replica in self.replicas.values_mut() {
343 replica.remove_collection(id);
344 }
345
346 self.collections.remove(&id);
348 }
349
350 fn add_replica_state(
351 &mut self,
352 id: ReplicaId,
353 client: ReplicaClient,
354 config: ReplicaConfig,
355 epoch: u64,
356 ) -> Result<(), read_holds::ReadHoldIssuerHungUp> {
357 let log_ids: BTreeSet<_> = config.logging.index_logs.values().copied().collect();
358
359 let metrics = self.metrics.for_replica(id);
360 let mut replica = ReplicaState::new(
361 id,
362 client,
363 config,
364 metrics,
365 self.introspection_tx.clone(),
366 epoch,
367 );
368
369 let mut shutdown_input = None;
371 for (collection_id, collection) in &self.collections {
372 if (collection.log_collection && !log_ids.contains(collection_id))
375 || collection.target_replica.is_some_and(|rid| rid != id)
376 {
377 continue;
378 }
379
380 let as_of = if collection.log_collection {
381 Antichain::from_elem(Timestamp::MIN)
386 } else {
387 collection.read_frontier().to_owned()
388 };
389
390 let mut input_read_holds = Vec::with_capacity(collection.storage_dependencies.len());
398 let mut hung_up = Vec::new();
399 for hold in collection.storage_dependencies.values() {
400 match hold.try_clone() {
401 Ok(hold) => input_read_holds.push(hold),
402 Err(read_holds::ReadHoldIssuerHungUp(input_id)) => hung_up.push(input_id),
403 }
404 }
405 if !hung_up.is_empty() {
406 tracing::error!(
407 replica_id = %id,
408 %collection_id,
409 ?hung_up,
410 "giving up on adding replica collections: storage read hold issuers hung \
411 up, the process is shutting down",
412 );
413 shutdown_input = hung_up.into_iter().next();
414 break;
415 }
416
417 replica.add_collection(*collection_id, as_of, input_read_holds);
418 }
419
420 self.replicas.insert(id, replica);
421
422 match shutdown_input {
423 Some(input_id) => Err(read_holds::ReadHoldIssuerHungUp(input_id)),
424 None => Ok(()),
425 }
426 }
427
428 fn deliver_response(&self, response: ComputeControllerResponse) {
430 let _ = self.response_tx.send(response);
433 }
434
435 fn deliver_introspection_updates(&self, type_: IntrospectionType, updates: Vec<(Row, Diff)>) {
437 let _ = self.introspection_tx.send((type_, updates));
440 }
441
442 fn replica_exists(&self, id: ReplicaId) -> bool {
444 self.replicas.contains_key(&id)
445 }
446
447 fn peeks_targeting(&self, replica_id: ReplicaId) -> impl Iterator<Item = (Uuid, &PendingPeek)> {
449 self.peeks.iter().filter_map(move |(uuid, peek)| {
450 if peek.target_replica == Some(replica_id) {
451 Some((*uuid, peek))
452 } else {
453 None
454 }
455 })
456 }
457
458 fn subscribes_targeting(&self, replica_id: ReplicaId) -> impl Iterator<Item = GlobalId> + '_ {
460 self.subscribes.keys().copied().filter(move |id| {
461 let collection = self.expect_collection(*id);
462 collection.target_replica == Some(replica_id)
463 })
464 }
465
466 fn update_frontier_introspection(&mut self) {
475 for collection in self.collections.values_mut() {
476 collection
477 .introspection
478 .observe_frontiers(&collection.read_frontier(), &collection.write_frontier());
479 }
480
481 for replica in self.replicas.values_mut() {
482 for collection in replica.collections.values_mut() {
483 collection
484 .introspection
485 .observe_frontier(&collection.write_frontier);
486 }
487 }
488 }
489
490 fn refresh_state_metrics(&self) {
499 let unscheduled_collections_count =
500 self.collections.values().filter(|c| !c.scheduled).count();
501 let connected_replica_count = self
502 .replicas
503 .values()
504 .filter(|r| r.client.is_connected())
505 .count();
506
507 self.metrics
508 .replica_count
509 .set(u64::cast_from(self.replicas.len()));
510 self.metrics
511 .collection_count
512 .set(u64::cast_from(self.collections.len()));
513 self.metrics
514 .collection_unscheduled_count
515 .set(u64::cast_from(unscheduled_collections_count));
516 self.metrics
517 .peek_count
518 .set(u64::cast_from(self.peeks.len()));
519 self.metrics
520 .subscribe_count
521 .set(u64::cast_from(self.subscribes.len()));
522 self.metrics
523 .copy_to_count
524 .set(u64::cast_from(self.copy_tos.len()));
525 self.metrics
526 .connected_replica_count
527 .set(u64::cast_from(connected_replica_count));
528 }
529
530 fn refresh_wallclock_lag(&mut self) {
549 let frontier_lag = |frontier: &Antichain<Timestamp>| match frontier.as_option() {
550 Some(ts) => (self.wallclock_lag)(ts.clone()),
551 None => Duration::ZERO,
552 };
553
554 let now_ms = (self.now)();
555 let histogram_period = WallclockLagHistogramPeriod::from_epoch_millis(now_ms, &self.dyncfg);
556 let histogram_labels = match &self.workload_class {
557 Some(wc) => [("workload_class", wc.clone())].into(),
558 None => BTreeMap::new(),
559 };
560
561 let readable_storage_collections: BTreeSet<_> = self
564 .collections
565 .keys()
566 .filter_map(|id| {
567 let frontiers = self.storage_collections.collection_frontiers(*id).ok()?;
568 PartialOrder::less_than(&frontiers.read_capabilities, &frontiers.write_frontier)
569 .then_some(*id)
570 })
571 .collect();
572
573 for (id, collection) in &mut self.collections {
575 let write_frontier = collection.write_frontier();
576 let readable = if self.storage_collections.check_exists(*id).is_ok() {
577 readable_storage_collections.contains(id)
578 } else {
579 PartialOrder::less_than(&collection.read_frontier(), &write_frontier)
580 };
581
582 if let Some(stash) = &mut collection.wallclock_lag_histogram_stash {
583 let bucket = if readable {
584 let lag = frontier_lag(&write_frontier);
585 let lag = lag.as_secs().next_power_of_two();
586 WallclockLag::Seconds(lag)
587 } else {
588 WallclockLag::Undefined
589 };
590
591 let key = (histogram_period, bucket, histogram_labels.clone());
592 *stash.entry(key).or_default() += Diff::ONE;
593 }
594 }
595
596 for replica in self.replicas.values_mut() {
598 for (id, collection) in &mut replica.collections {
599 let readable = readable_storage_collections.contains(id) || collection.hydrated();
604
605 let lag = if readable {
606 let lag = frontier_lag(&collection.write_frontier);
607 WallclockLag::Seconds(lag.as_secs())
608 } else {
609 WallclockLag::Undefined
610 };
611
612 if let Some(wallclock_lag_max) = &mut collection.wallclock_lag_max {
613 *wallclock_lag_max = (*wallclock_lag_max).max(lag);
614 }
615
616 if let Some(metrics) = &mut collection.metrics {
617 let secs = lag.unwrap_seconds_or(u64::MAX);
620 metrics.wallclock_lag.observe(secs);
621 };
622 }
623 }
624
625 self.maybe_record_wallclock_lag();
627 }
628
629 fn maybe_record_wallclock_lag(&mut self) {
637 if self.read_only {
638 return;
639 }
640
641 let duration_trunc = |datetime: DateTime<_>, interval| {
642 let td = TimeDelta::from_std(interval).ok()?;
643 datetime.duration_trunc(td).ok()
644 };
645
646 let interval = WALLCLOCK_LAG_RECORDING_INTERVAL.get(&self.dyncfg);
647 let now_dt = mz_ore::now::to_datetime((self.now)());
648 let now_trunc = duration_trunc(now_dt, interval).unwrap_or_else(|| {
649 soft_panic_or_log!("excessive wallclock lag recording interval: {interval:?}");
650 let default = WALLCLOCK_LAG_RECORDING_INTERVAL.default();
651 duration_trunc(now_dt, *default).unwrap()
652 });
653 if now_trunc <= self.wallclock_lag_last_recorded {
654 return;
655 }
656
657 let now_ts: CheckedTimestamp<_> = now_trunc.try_into().expect("must fit");
658
659 let mut history_updates = Vec::new();
660 for (replica_id, replica) in &mut self.replicas {
661 for (collection_id, collection) in &mut replica.collections {
662 let Some(wallclock_lag_max) = &mut collection.wallclock_lag_max else {
663 continue;
664 };
665
666 let max_lag = std::mem::replace(wallclock_lag_max, WallclockLag::MIN);
667 let row = Row::pack_slice(&[
668 Datum::String(&collection_id.to_string()),
669 Datum::String(&replica_id.to_string()),
670 max_lag.into_interval_datum(),
671 Datum::TimestampTz(now_ts),
672 ]);
673 history_updates.push((row, Diff::ONE));
674 }
675 }
676 if !history_updates.is_empty() {
677 self.deliver_introspection_updates(
678 IntrospectionType::WallclockLagHistory,
679 history_updates,
680 );
681 }
682
683 let mut histogram_updates = Vec::new();
684 let mut row_buf = Row::default();
685 for (collection_id, collection) in &mut self.collections {
686 let Some(stash) = &mut collection.wallclock_lag_histogram_stash else {
687 continue;
688 };
689
690 for ((period, lag, labels), count) in std::mem::take(stash) {
691 let mut packer = row_buf.packer();
692 packer.extend([
693 Datum::TimestampTz(period.start),
694 Datum::TimestampTz(period.end),
695 Datum::String(&collection_id.to_string()),
696 lag.into_uint64_datum(),
697 ]);
698 let labels = labels.iter().map(|(k, v)| (*k, Datum::String(v)));
699 packer.push_dict(labels);
700
701 histogram_updates.push((row_buf.clone(), count));
702 }
703 }
704 if !histogram_updates.is_empty() {
705 self.deliver_introspection_updates(
706 IntrospectionType::WallclockLagHistogram,
707 histogram_updates,
708 );
709 }
710
711 self.wallclock_lag_last_recorded = now_trunc;
712 }
713
714 #[mz_ore::instrument(level = "debug")]
720 pub fn collection_hydrated(&self, collection_id: GlobalId) -> Result<bool, CollectionMissing> {
721 let mut hosting_replicas = self.replicas_hosting(collection_id)?.peekable();
722 if hosting_replicas.peek().is_none() {
723 return Ok(true);
724 }
725 for replica_state in hosting_replicas {
726 if replica_state.expect_collection(collection_id).hydrated() {
727 return Ok(true);
728 }
729 }
730
731 Ok(false)
732 }
733
734 #[mz_ore::instrument(level = "debug")]
750 pub fn collections_ready_on_replicas(
751 &self,
752 target_replica_ids: Option<Vec<ReplicaId>>,
753 exclude_collections: &BTreeSet<GlobalId>,
754 allowed_lag: Option<Timestamp>,
755 reference_replica_ids: &BTreeSet<ReplicaId>,
756 ) -> Result<bool, HydrationCheckBadTarget> {
757 if self.replicas.is_empty() {
758 return Ok(true);
759 }
760 let target_replicas: BTreeSet<ReplicaId> = self
761 .replicas
762 .keys()
763 .filter_map(|id| match target_replica_ids {
764 None => Some(id.clone()),
765 Some(ref ids) if ids.contains(id) => Some(id.clone()),
766 Some(_) => None,
767 })
768 .collect();
769 if let Some(targets) = target_replica_ids {
770 if target_replicas.is_empty() {
771 return Err(HydrationCheckBadTarget(targets));
772 }
773 }
774
775 let mut unhydrated = BTreeSet::new();
776 let mut lagging_ticks = BTreeMap::new();
777 let mut awaiting_completion = BTreeSet::new();
778 for (id, _collection) in self.collections_iter() {
779 if id.is_transient() || exclude_collections.contains(&id) {
780 continue;
781 }
782
783 let replicas = self
786 .replicas_hosting(id)
787 .expect("collection must exist")
788 .map(|replica| (replica.id, replica.expect_collection(id)));
789
790 match classify_collection_readiness(
791 replicas,
792 &target_replicas,
793 reference_replica_ids,
794 allowed_lag,
795 ) {
796 CollectionReadiness::Ready => {}
797 CollectionReadiness::Lagging { lag: Some(lag) } => {
801 lagging_ticks.insert(id, lag);
802 }
803 CollectionReadiness::Lagging { lag: None } => {
804 awaiting_completion.insert(id);
805 }
806 CollectionReadiness::Unhydrated => {
807 unhydrated.insert(id);
808 }
809 }
810 }
811
812 let ready =
813 unhydrated.is_empty() && lagging_ticks.is_empty() && awaiting_completion.is_empty();
814 if !ready {
815 tracing::info!(
819 replicas = ?target_replicas,
820 reference = ?reference_replica_ids,
821 unhydrated = ?unhydrated,
822 ?lagging_ticks,
823 ?awaiting_completion,
824 ?allowed_lag,
825 "collections are not ready on any target replica",
826 );
827 }
828
829 Ok(ready)
830 }
831
832 fn cleanup_collections(&mut self) {
848 let to_remove: Vec<_> = self
849 .collections_iter()
850 .filter(|(id, collection)| {
851 collection.dropped
852 && collection.shared.lock_read_capabilities(|c| c.is_empty())
853 && self
854 .replicas
855 .values()
856 .all(|r| r.collection_frontiers_empty(*id))
857 })
858 .map(|(id, _collection)| id)
859 .collect();
860
861 for id in to_remove {
862 self.remove_collection(id);
863 }
864 }
865
866 #[mz_ore::instrument(level = "debug")]
870 pub fn dump(&self) -> Result<serde_json::Value, anyhow::Error> {
871 let Self {
878 build_info: _,
879 storage_collections: _,
880 peek_stash_persist_location: _,
881 initialized,
882 read_only,
883 workload_class,
884 replicas,
885 replica_dyncfg_overrides: _,
886 collections,
887 log_sources: _,
888 peeks,
889 subscribes,
890 copy_tos,
891 history: _,
892 command_rx: _,
893 response_tx: _,
894 introspection_tx: _,
895 metrics: _,
896 dyncfg: _,
897 now: _,
898 wallclock_lag: _,
899 wallclock_lag_last_recorded,
900 read_hold_tx: _,
901 replica_tx: _,
902 replica_rx: _,
903 } = self;
904
905 let replicas: BTreeMap<_, _> = replicas
906 .iter()
907 .map(|(id, replica)| Ok((id.to_string(), replica.dump()?)))
908 .collect::<Result<_, anyhow::Error>>()?;
909 let collections: BTreeMap<_, _> = collections
910 .iter()
911 .map(|(id, collection)| (id.to_string(), format!("{collection:?}")))
912 .collect();
913 let peeks: BTreeMap<_, _> = peeks
914 .iter()
915 .map(|(uuid, peek)| (uuid.to_string(), format!("{peek:?}")))
916 .collect();
917 let subscribes: BTreeMap<_, _> = subscribes
918 .iter()
919 .map(|(id, subscribe)| (id.to_string(), format!("{subscribe:?}")))
920 .collect();
921 let copy_tos: Vec<_> = copy_tos.iter().map(|id| id.to_string()).collect();
922 let wallclock_lag_last_recorded = format!("{wallclock_lag_last_recorded:?}");
923
924 Ok(serde_json::json!({
925 "initialized": initialized,
926 "read_only": read_only,
927 "workload_class": workload_class,
928 "replicas": replicas,
929 "collections": collections,
930 "peeks": peeks,
931 "subscribes": subscribes,
932 "copy_tos": copy_tos,
933 "wallclock_lag_last_recorded": wallclock_lag_last_recorded,
934 }))
935 }
936
937 pub(super) fn collection_write_frontier(
939 &self,
940 id: GlobalId,
941 ) -> Result<Antichain<Timestamp>, CollectionMissing> {
942 Ok(self.collection(id)?.write_frontier())
943 }
944}
945
946impl Instance {
947 pub(super) fn new(
948 build_info: &'static BuildInfo,
949 storage: StorageCollections,
950 peek_stash_persist_location: PersistLocation,
951 arranged_logs: Vec<(LogVariant, GlobalId, SharedCollectionState)>,
952 metrics: InstanceMetrics,
953 now: NowFn,
954 wallclock_lag: WallclockLagFn<Timestamp>,
955 dyncfg: Arc<ConfigSet>,
956 command_rx: mpsc::UnboundedReceiver<Command>,
957 response_tx: mpsc::UnboundedSender<ComputeControllerResponse>,
958 read_hold_tx: read_holds::ChangeTx,
959 introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
960 read_only: bool,
961 ) -> Self {
962 let mut collections = BTreeMap::new();
963 let mut log_sources = BTreeMap::new();
964 for (log, id, shared) in arranged_logs {
965 let collection = CollectionState::new_log_collection(
966 id,
967 shared,
968 Arc::clone(&read_hold_tx),
969 introspection_tx.clone(),
970 );
971 collections.insert(id, collection);
972 log_sources.insert(log, id);
973 }
974
975 let history = ComputeCommandHistory::new(metrics.for_history());
976
977 let send_count = metrics.response_send_count.clone();
978 let recv_count = metrics.response_recv_count.clone();
979 let (replica_tx, replica_rx) = instrumented_unbounded_channel(send_count, recv_count);
980
981 let now_dt = mz_ore::now::to_datetime(now());
982
983 Self {
984 build_info,
985 storage_collections: storage,
986 peek_stash_persist_location,
987 initialized: false,
988 read_only,
989 workload_class: None,
990 replicas: Default::default(),
991 replica_dyncfg_overrides: Default::default(),
992 collections,
993 log_sources,
994 peeks: Default::default(),
995 subscribes: Default::default(),
996 copy_tos: Default::default(),
997 history,
998 command_rx,
999 response_tx,
1000 introspection_tx,
1001 metrics,
1002 dyncfg,
1003 now,
1004 wallclock_lag,
1005 wallclock_lag_last_recorded: now_dt,
1006 read_hold_tx,
1007 replica_tx,
1008 replica_rx,
1009 }
1010 }
1011
1012 pub(super) async fn run(mut self) {
1013 self.send(ComputeCommand::Hello {
1014 nonce: Uuid::default(),
1017 });
1018
1019 let instance_config = InstanceConfig {
1020 peek_stash_persist_location: self.peek_stash_persist_location.clone(),
1021 logging: Default::default(),
1025 expiration_offset: Default::default(),
1026 arrangement_dictionary_compression: Default::default(),
1027 initial_config: Default::default(),
1028 };
1029
1030 self.send(ComputeCommand::CreateInstance(Box::new(instance_config)));
1031
1032 loop {
1033 tokio::select! {
1034 command = self.command_rx.recv() => match command {
1035 Some(cmd) => cmd(&mut self),
1036 None => break,
1037 },
1038 response = self.replica_rx.recv() => match response {
1039 Some(response) => self.handle_response(response),
1040 None => unreachable!("self owns a sender side of the channel"),
1041 }
1042 }
1043 }
1044 }
1045
1046 #[mz_ore::instrument(level = "debug")]
1048 pub fn update_configuration(&mut self, config_params: ComputeParameters) {
1049 if let Some(workload_class) = &config_params.workload_class {
1050 self.workload_class = workload_class.clone();
1051 }
1052
1053 let command = ComputeCommand::UpdateConfiguration(Box::new(config_params));
1054 self.send(command);
1055 }
1056
1057 #[mz_ore::instrument(level = "debug")]
1062 pub fn initialization_complete(&mut self) {
1063 if !self.initialized {
1065 self.send(ComputeCommand::InitializationComplete);
1066 self.initialized = true;
1067 }
1068 }
1069
1070 #[mz_ore::instrument(level = "debug")]
1074 pub fn allow_writes(&mut self, collection_id: GlobalId) -> Result<(), CollectionMissing> {
1075 let collection = self.collection_mut(collection_id)?;
1076
1077 if !collection.read_only {
1079 return Ok(());
1080 }
1081
1082 let as_of = collection.read_frontier();
1084
1085 if as_of.is_empty() {
1088 return Ok(());
1089 }
1090
1091 collection.read_only = false;
1092 self.send(ComputeCommand::AllowWrites(collection_id));
1093
1094 Ok(())
1095 }
1096
1097 #[mz_ore::instrument(level = "debug")]
1107 pub fn shutdown(&mut self) {
1108 let (_tx, rx) = mpsc::unbounded_channel();
1110 self.command_rx = rx;
1111
1112 let stray_replicas: Vec<_> = self.replicas.keys().collect();
1113 soft_assert_or_log!(
1114 stray_replicas.is_empty(),
1115 "dropped instance still has provisioned replicas: {stray_replicas:?}",
1116 );
1117 }
1118
1119 fn initiate_shutdown(&mut self) {
1125 let (_tx, rx) = mpsc::unbounded_channel();
1128 self.command_rx = rx;
1129 }
1130
1131 #[mz_ore::instrument(level = "debug")]
1133 fn send(&mut self, cmd: ComputeCommand) {
1134 self.history.push(cmd.clone());
1139
1140 let target_replica = self.target_replica(&cmd);
1141
1142 let overrides = &self.replica_dyncfg_overrides;
1145 let dyncfg = &self.dyncfg;
1146
1147 if let Some(rid) = target_replica {
1148 if let Some(replica) = self.replicas.get_mut(&rid) {
1149 let cmd = Self::specialize_command_for_replica(cmd, rid, overrides, dyncfg);
1150 let _ = replica.client.send(cmd);
1151 }
1152 } else {
1153 for (rid, replica) in self.replicas.iter_mut() {
1154 let cmd =
1155 Self::specialize_command_for_replica(cmd.clone(), *rid, overrides, dyncfg);
1156 let _ = replica.client.send(cmd);
1157 }
1158 }
1159 }
1160
1161 fn specialize_command_for_replica(
1169 mut cmd: ComputeCommand,
1170 replica_id: ReplicaId,
1171 overrides: &BTreeMap<ReplicaId, ConfigUpdates>,
1172 dyncfg: &ConfigSet,
1173 ) -> ComputeCommand {
1174 let over = overrides.get(&replica_id);
1175 match &mut cmd {
1176 ComputeCommand::UpdateConfiguration(params) => {
1177 if let Some(over) = over
1178 && !over.updates.is_empty()
1179 {
1180 params.dyncfg_updates.extend(over.clone());
1181 }
1182 }
1183 ComputeCommand::CreateInstance(config) => {
1184 let mut initial = ConfigUpdates::from(dyncfg);
1185 if let Some(over) = over {
1186 initial.extend(over.clone());
1187 }
1188 config.initial_config = initial;
1189 }
1190 _ => {}
1191 }
1192 cmd
1193 }
1194
1195 pub(super) fn update_replica_dyncfg_overrides(
1199 &mut self,
1200 overrides: BTreeMap<ReplicaId, ConfigUpdates>,
1201 ) {
1202 self.replica_dyncfg_overrides = overrides;
1203 }
1204
1205 fn target_replica(&self, cmd: &ComputeCommand) -> Option<ReplicaId> {
1213 match &cmd {
1214 ComputeCommand::Schedule(id)
1215 | ComputeCommand::AllowWrites(id)
1216 | ComputeCommand::AllowCompaction { id, .. } => {
1217 self.expect_collection(*id).target_replica
1218 }
1219 ComputeCommand::CreateDataflow(desc) => {
1220 let mut target_replica = None;
1221 for id in desc.export_ids() {
1222 if let Some(replica) = self.expect_collection(id).target_replica {
1223 if target_replica.is_some() {
1224 assert_eq!(target_replica, Some(replica));
1225 }
1226 target_replica = Some(replica);
1227 }
1228 }
1229 target_replica
1230 }
1231 ComputeCommand::Peek(_)
1233 | ComputeCommand::Hello { .. }
1234 | ComputeCommand::CreateInstance(_)
1235 | ComputeCommand::InitializationComplete
1236 | ComputeCommand::UpdateConfiguration(_)
1237 | ComputeCommand::CancelPeek { .. } => None,
1238 }
1239 }
1240
1241 #[mz_ore::instrument(level = "debug")]
1243 pub fn add_replica(
1244 &mut self,
1245 id: ReplicaId,
1246 mut config: ReplicaConfig,
1247 epoch: Option<u64>,
1248 ) -> Result<(), ReplicaExists> {
1249 if self.replica_exists(id) {
1250 return Err(ReplicaExists(id));
1251 }
1252
1253 config.logging.index_logs = self.log_sources.clone();
1254
1255 let epoch = epoch.unwrap_or(1);
1256 let metrics = self.metrics.for_replica(id);
1257 let client = ReplicaClient::spawn(
1258 id,
1259 self.build_info,
1260 config.clone(),
1261 epoch,
1262 metrics.clone(),
1263 Arc::clone(&self.dyncfg),
1264 self.replica_tx.clone(),
1265 );
1266
1267 self.history.reduce();
1269
1270 self.history.update_source_uppers(&self.storage_collections);
1272
1273 for command in self.history.iter() {
1275 if let Some(target_replica) = self.target_replica(command)
1277 && target_replica != id
1278 {
1279 continue;
1280 }
1281
1282 let command = Self::specialize_command_for_replica(
1285 command.clone(),
1286 id,
1287 &self.replica_dyncfg_overrides,
1288 &self.dyncfg,
1289 );
1290 if client.send(command).is_err() {
1291 tracing::warn!("Replica {:?} connection terminated during hydration", id);
1294 break;
1295 }
1296 }
1297
1298 if self.add_replica_state(id, client, config, epoch).is_err() {
1300 self.initiate_shutdown();
1306 }
1307
1308 Ok(())
1309 }
1310
1311 #[mz_ore::instrument(level = "debug")]
1313 pub fn remove_replica(&mut self, id: ReplicaId) -> Result<(), ReplicaMissing> {
1314 let replica = self.replicas.remove(&id).ok_or(ReplicaMissing(id))?;
1315
1316 self.replica_dyncfg_overrides.remove(&id);
1320
1321 for (collection_id, replica_collection) in &replica.collections {
1329 let collection = self.collections.get(collection_id);
1330 for replica_hold in &replica_collection.input_read_holds {
1331 let input_id = replica_hold.id();
1332 let global_hold = collection.and_then(|c| c.storage_dependencies.get(&input_id));
1333 let unprotected = global_hold
1334 .is_none_or(|h| PartialOrder::less_than(replica_hold.since(), h.since()));
1335 if unprotected {
1336 tracing::warn!(
1337 replica_id = %id,
1338 %collection_id,
1339 %input_id,
1340 replica_hold_since = ?replica_hold.since(),
1341 global_hold_since = ?global_hold.map(|h| h.since()),
1342 "dropping per-replica read hold without equivalent global read hold",
1343 );
1344 }
1345 }
1346 }
1347 drop(replica);
1348
1349 let to_drop: Vec<_> = self.subscribes_targeting(id).collect();
1353 for subscribe_id in to_drop {
1354 let subscribe = self.subscribes.remove(&subscribe_id).unwrap();
1355 let response = ComputeControllerResponse::SubscribeResponse(
1356 subscribe_id,
1357 SubscribeBatch {
1358 lower: subscribe.frontier.clone(),
1359 upper: subscribe.frontier,
1360 updates: Err(ERROR_TARGET_REPLICA_FAILED.into()),
1361 },
1362 );
1363 self.deliver_response(response);
1364 }
1365
1366 let mut peek_responses = Vec::new();
1371 let mut to_drop = Vec::new();
1372 for (uuid, peek) in self.peeks_targeting(id) {
1373 peek_responses.push(ComputeControllerResponse::PeekNotification(
1374 uuid,
1375 PeekNotification::Error(ERROR_TARGET_REPLICA_FAILED.into()),
1376 peek.otel_ctx.clone(),
1377 ));
1378 to_drop.push(uuid);
1379 }
1380 for response in peek_responses {
1381 self.deliver_response(response);
1382 }
1383 for uuid in to_drop {
1384 let response =
1385 PeekResponse::Error(ProtocolPeekError::unstructured(ERROR_TARGET_REPLICA_FAILED));
1386 self.finish_peek(uuid, response);
1387 }
1388
1389 self.forward_implied_capabilities();
1392
1393 Ok(())
1394 }
1395
1396 fn rehydrate_replica(&mut self, id: ReplicaId) {
1402 let config = self.replicas[&id].config.clone();
1403 let epoch = self.replicas[&id].epoch + 1;
1404
1405 self.remove_replica(id).expect("replica must exist");
1406 let result = self.add_replica(id, config, Some(epoch));
1407
1408 match result {
1409 Ok(()) => (),
1410 Err(ReplicaExists(_)) => unreachable!("replica was removed"),
1411 }
1412 }
1413
1414 fn rehydrate_failed_replicas(&mut self) {
1416 let replicas = self.replicas.iter();
1417 let failed_replicas: Vec<_> = replicas
1418 .filter_map(|(id, replica)| replica.client.is_failed().then_some(*id))
1419 .collect();
1420
1421 for replica_id in failed_replicas {
1422 self.rehydrate_replica(replica_id);
1423 }
1424 }
1425
1426 #[mz_ore::instrument(level = "debug")]
1431 pub fn create_dataflow(
1432 &mut self,
1433 dataflow: DataflowDescription<mz_compute_types::plan::LirRelationExpr, ()>,
1434 import_read_holds: Vec<ReadHold>,
1435 mut shared_collection_state: BTreeMap<GlobalId, SharedCollectionState>,
1436 target_replica: Option<ReplicaId>,
1437 ) -> Result<(), DataflowCreationError> {
1438 use DataflowCreationError::*;
1439
1440 if let Some(replica_id) = target_replica {
1444 if !self.replica_exists(replica_id) {
1445 return Err(ReplicaMissing(replica_id));
1446 }
1447 }
1448
1449 let as_of = dataflow.as_of.as_ref().ok_or(MissingAsOf)?;
1451 if as_of.is_empty() && dataflow.subscribe_ids().next().is_some() {
1452 return Err(EmptyAsOfForSubscribe);
1453 }
1454 if as_of.is_empty() && dataflow.copy_to_ids().next().is_some() {
1455 return Err(EmptyAsOfForCopyTo);
1456 }
1457
1458 let mut storage_dependencies = BTreeMap::new();
1460 let mut compute_dependencies = BTreeMap::new();
1461
1462 let mut replica_input_read_holds = Vec::new();
1467
1468 let mut import_read_holds: BTreeMap<_, _> =
1469 import_read_holds.into_iter().map(|r| (r.id(), r)).collect();
1470
1471 for &id in dataflow.source_imports.keys() {
1472 let mut read_hold = import_read_holds.remove(&id).ok_or(ReadHoldMissing(id))?;
1473 replica_input_read_holds.push(read_hold.clone());
1474
1475 read_hold
1476 .try_downgrade(as_of.clone())
1477 .map_err(|_| ReadHoldInsufficient(id))?;
1478 storage_dependencies.insert(id, read_hold);
1479 }
1480
1481 for &id in dataflow.index_imports.keys() {
1482 let mut read_hold = import_read_holds.remove(&id).ok_or(ReadHoldMissing(id))?;
1483 read_hold
1484 .try_downgrade(as_of.clone())
1485 .map_err(|_| ReadHoldInsufficient(id))?;
1486 compute_dependencies.insert(id, read_hold);
1487 }
1488
1489 if as_of.is_empty() {
1492 replica_input_read_holds = Default::default();
1493 }
1494
1495 for export_id in dataflow.export_ids() {
1497 let shared = shared_collection_state
1498 .remove(&export_id)
1499 .unwrap_or_else(|| SharedCollectionState::new(as_of.clone()));
1500 let write_only = dataflow.sink_exports.contains_key(&export_id);
1501 let storage_sink = dataflow.persist_sink_ids().any(|id| id == export_id);
1502
1503 self.add_collection(
1504 export_id,
1505 as_of.clone(),
1506 shared,
1507 storage_dependencies.clone(),
1508 compute_dependencies.clone(),
1509 replica_input_read_holds.clone(),
1510 write_only,
1511 storage_sink,
1512 dataflow.initial_storage_as_of.clone(),
1513 dataflow.refresh_schedule.clone(),
1514 target_replica,
1515 );
1516
1517 if let Ok(frontiers) = self.storage_collections.collection_frontiers(export_id) {
1520 self.maybe_update_global_write_frontier(export_id, frontiers.write_frontier);
1521 }
1522 }
1523
1524 for subscribe_id in dataflow.subscribe_ids() {
1526 self.subscribes
1527 .insert(subscribe_id, ActiveSubscribe::default());
1528 }
1529
1530 for copy_to_id in dataflow.copy_to_ids() {
1532 self.copy_tos.insert(copy_to_id);
1533 }
1534
1535 let mut source_imports = BTreeMap::new();
1538 for (id, import) in dataflow.source_imports {
1539 let frontiers = self
1540 .storage_collections
1541 .collection_frontiers(id)
1542 .expect("collection exists");
1543
1544 let collection_metadata = self
1545 .storage_collections
1546 .collection_metadata(id)
1547 .expect("we have a read hold on this collection");
1548
1549 let desc = SourceInstanceDesc {
1550 storage_metadata: collection_metadata.clone(),
1551 arguments: import.desc.arguments,
1552 typ: import.desc.typ.clone(),
1553 };
1554 source_imports.insert(
1555 id,
1556 mz_compute_types::dataflows::SourceImport {
1557 desc,
1558 monotonic: import.monotonic,
1559 with_snapshot: import.with_snapshot,
1560 upper: frontiers.write_frontier,
1561 },
1562 );
1563 }
1564
1565 let mut sink_exports = BTreeMap::new();
1566 for (id, se) in dataflow.sink_exports {
1567 let connection = match se.connection {
1568 ComputeSinkConnection::MaterializedView(conn) => {
1569 let metadata = self
1570 .storage_collections
1571 .collection_metadata(id)
1572 .map_err(|_| CollectionMissing(id))?
1573 .clone();
1574 let conn = MaterializedViewSinkConnection {
1575 value_desc: conn.value_desc,
1576 storage_metadata: metadata,
1577 };
1578 ComputeSinkConnection::MaterializedView(conn)
1579 }
1580 ComputeSinkConnection::Subscribe(conn) => ComputeSinkConnection::Subscribe(conn),
1581 ComputeSinkConnection::CopyToS3Oneshot(conn) => {
1582 ComputeSinkConnection::CopyToS3Oneshot(conn)
1583 }
1584 ComputeSinkConnection::MetricSink(conn) => ComputeSinkConnection::MetricSink(conn),
1585 };
1586 let desc = ComputeSinkDesc {
1587 from: se.from,
1588 from_desc: se.from_desc,
1589 connection,
1590 with_snapshot: se.with_snapshot,
1591 up_to: se.up_to,
1592 non_null_assertions: se.non_null_assertions,
1593 refresh_schedule: se.refresh_schedule,
1594 };
1595 sink_exports.insert(id, desc);
1596 }
1597
1598 let objects_to_build = dataflow
1600 .objects_to_build
1601 .into_iter()
1602 .map(|object| BuildDesc {
1603 id: object.id,
1604 plan: RenderPlan::try_from(object.plan).expect("valid plan"),
1605 })
1606 .collect();
1607
1608 let augmented_dataflow = DataflowDescription {
1609 source_imports,
1610 sink_exports,
1611 objects_to_build,
1612 index_imports: dataflow.index_imports,
1614 index_exports: dataflow.index_exports,
1615 as_of: dataflow.as_of.clone(),
1616 until: dataflow.until,
1617 initial_storage_as_of: dataflow.initial_storage_as_of,
1618 refresh_schedule: dataflow.refresh_schedule,
1619 debug_name: dataflow.debug_name,
1620 time_dependence: dataflow.time_dependence,
1621 };
1622
1623 if augmented_dataflow.is_transient() {
1624 tracing::debug!(
1625 name = %augmented_dataflow.debug_name,
1626 import_ids = %augmented_dataflow.display_import_ids(),
1627 export_ids = %augmented_dataflow.display_export_ids(),
1628 as_of = ?augmented_dataflow.as_of.as_ref().unwrap().elements(),
1629 until = ?augmented_dataflow.until.elements(),
1630 "creating dataflow",
1631 );
1632 } else {
1633 tracing::info!(
1634 name = %augmented_dataflow.debug_name,
1635 import_ids = %augmented_dataflow.display_import_ids(),
1636 export_ids = %augmented_dataflow.display_export_ids(),
1637 as_of = ?augmented_dataflow.as_of.as_ref().unwrap().elements(),
1638 until = ?augmented_dataflow.until.elements(),
1639 "creating dataflow",
1640 );
1641 }
1642
1643 if as_of.is_empty() {
1646 tracing::info!(
1647 name = %augmented_dataflow.debug_name,
1648 "not sending `CreateDataflow`, because of empty `as_of`",
1649 );
1650 } else {
1651 let collections: Vec<_> = augmented_dataflow.export_ids().collect();
1652 self.send(ComputeCommand::CreateDataflow(Box::new(augmented_dataflow)));
1653
1654 for id in collections {
1655 self.maybe_schedule_collection(id);
1656 }
1657 }
1658
1659 Ok(())
1660 }
1661
1662 fn maybe_schedule_collection(&mut self, id: GlobalId) {
1668 let collection = self.expect_collection(id);
1669
1670 if collection.scheduled {
1672 return;
1673 }
1674
1675 let as_of = collection.read_frontier();
1676
1677 if as_of.is_empty() {
1680 return;
1681 }
1682
1683 let ready = if id.is_transient() {
1684 true
1690 } else {
1691 let not_self_dep = |x: &GlobalId| *x != id;
1697
1698 let mut deps_scheduled = true;
1701
1702 let compute_deps = collection.compute_dependency_ids().filter(not_self_dep);
1707 let mut compute_frontiers = Vec::new();
1708 for id in compute_deps {
1709 let dep = &self.expect_collection(id);
1710 deps_scheduled &= dep.scheduled;
1711 compute_frontiers.push(dep.write_frontier());
1712 }
1713
1714 let storage_deps = collection.storage_dependency_ids().filter(not_self_dep);
1715 let storage_frontiers = self
1716 .storage_collections
1717 .collections_frontiers(storage_deps.collect())
1718 .expect("must exist");
1719 let storage_frontiers = storage_frontiers.into_iter().map(|f| f.write_frontier);
1720
1721 let mut frontiers = compute_frontiers.into_iter().chain(storage_frontiers);
1722 let frontiers_ready =
1723 frontiers.all(|frontier| PartialOrder::less_than(&as_of, &frontier));
1724
1725 deps_scheduled && frontiers_ready
1726 };
1727
1728 if ready {
1729 self.send(ComputeCommand::Schedule(id));
1730 let collection = self.expect_collection_mut(id);
1731 collection.scheduled = true;
1732 }
1733 }
1734
1735 fn schedule_collections(&mut self) {
1737 let ids: Vec<_> = self.collections.keys().copied().collect();
1738 for id in ids {
1739 self.maybe_schedule_collection(id);
1740 }
1741 }
1742
1743 #[mz_ore::instrument(level = "debug")]
1746 pub fn drop_collections(&mut self, ids: Vec<GlobalId>) -> Result<(), CollectionMissing> {
1747 for id in &ids {
1748 let collection = self.collection_mut(*id)?;
1749
1750 collection.dropped = true;
1752
1753 collection.implied_read_hold.release();
1756 collection.warmup_read_hold.release();
1757
1758 self.subscribes.remove(id);
1761 self.copy_tos.remove(id);
1764 }
1765
1766 Ok(())
1767 }
1768
1769 #[mz_ore::instrument(level = "debug")]
1773 pub fn peek(
1774 &mut self,
1775 peek_target: PeekTarget,
1776 literal_constraints: Option<Vec<Row>>,
1777 uuid: Uuid,
1778 timestamp: Timestamp,
1779 result_desc: RelationDesc,
1780 finishing: RowSetFinishing,
1781 map_filter_project: mz_expr::SafeMfpPlan,
1782 mut read_hold: ReadHold,
1783 target_replica: Option<ReplicaId>,
1784 peek_response_tx: oneshot::Sender<PeekResponse>,
1785 ) -> Result<(), PeekError> {
1786 use PeekError::*;
1787
1788 let target_id = peek_target.id();
1789
1790 if read_hold.id() != target_id {
1792 return Err(ReadHoldIdMismatch(read_hold.id()));
1793 }
1794 read_hold
1795 .try_downgrade(Antichain::from_elem(timestamp.clone()))
1796 .map_err(|_| ReadHoldInsufficient(target_id))?;
1797
1798 if let Some(target) = target_replica {
1799 if !self.replica_exists(target) {
1800 return Err(ReplicaMissing(target));
1801 }
1802 }
1803
1804 let otel_ctx = OpenTelemetryContext::obtain();
1805
1806 self.peeks.insert(
1807 uuid,
1808 PendingPeek {
1809 target_replica,
1810 otel_ctx: otel_ctx.clone(),
1812 requested_at: Instant::now(),
1813 read_hold,
1814 peek_response_tx,
1815 limit: finishing.limit.map(usize::cast_from),
1816 offset: finishing.offset,
1817 },
1818 );
1819
1820 let peek = Peek {
1821 literal_constraints,
1822 uuid,
1823 timestamp,
1824 finishing,
1825 map_filter_project,
1826 otel_ctx,
1829 target: peek_target,
1830 result_desc,
1831 };
1832 self.send(ComputeCommand::Peek(Box::new(peek)));
1833
1834 Ok(())
1835 }
1836
1837 #[mz_ore::instrument(level = "debug")]
1839 pub fn cancel_peek(&mut self, uuid: Uuid, reason: PeekResponse) {
1840 let Some(peek) = self.peeks.get_mut(&uuid) else {
1841 tracing::warn!("did not find pending peek for {uuid}");
1842 return;
1843 };
1844
1845 let duration = peek.requested_at.elapsed();
1846 self.metrics
1847 .observe_peek_response(&PeekResponse::Canceled, duration);
1848
1849 let otel_ctx = peek.otel_ctx.clone();
1851 otel_ctx.attach_as_parent();
1852
1853 self.deliver_response(ComputeControllerResponse::PeekNotification(
1854 uuid,
1855 PeekNotification::Canceled,
1856 otel_ctx,
1857 ));
1858
1859 self.finish_peek(uuid, reason);
1862 }
1863
1864 #[mz_ore::instrument(level = "debug")]
1876 pub fn set_read_policy(
1877 &mut self,
1878 policies: Vec<(GlobalId, ReadPolicy)>,
1879 ) -> Result<(), ReadPolicyError> {
1880 for (id, _policy) in &policies {
1883 let collection = self.collection(*id)?;
1884 if collection.read_policy.is_none() {
1885 return Err(ReadPolicyError::WriteOnlyCollection(*id));
1886 }
1887 }
1888
1889 for (id, new_policy) in policies {
1890 let collection = self.expect_collection_mut(id);
1891 let new_since = new_policy.frontier(collection.write_frontier().borrow());
1892 let _ = collection.implied_read_hold.try_downgrade(new_since);
1893 collection.read_policy = Some(new_policy);
1894 }
1895
1896 Ok(())
1897 }
1898
1899 #[mz_ore::instrument(level = "debug")]
1907 fn maybe_update_global_write_frontier(
1908 &mut self,
1909 id: GlobalId,
1910 new_frontier: Antichain<Timestamp>,
1911 ) {
1912 let collection = self.expect_collection_mut(id);
1913
1914 let advanced = collection.shared.lock_write_frontier(|f| {
1915 let advanced = PartialOrder::less_than(f, &new_frontier);
1916 if advanced {
1917 f.clone_from(&new_frontier);
1918 }
1919 advanced
1920 });
1921
1922 if !advanced {
1923 return;
1924 }
1925
1926 let new_since = match &collection.read_policy {
1928 Some(read_policy) => {
1929 read_policy.frontier(new_frontier.borrow())
1932 }
1933 None => {
1934 Antichain::from_iter(
1943 new_frontier
1944 .iter()
1945 .map(|t| t.step_back().unwrap_or(Timestamp::MIN)),
1946 )
1947 }
1948 };
1949 let _ = collection.implied_read_hold.try_downgrade(new_since);
1950
1951 self.deliver_response(ComputeControllerResponse::FrontierUpper {
1953 id,
1954 upper: new_frontier,
1955 });
1956 }
1957
1958 pub(super) fn apply_read_hold_change(
1960 &mut self,
1961 id: GlobalId,
1962 mut update: ChangeBatch<Timestamp>,
1963 ) {
1964 let Some(collection) = self.collections.get_mut(&id) else {
1965 soft_panic_or_log!(
1966 "read hold change for absent collection (id={id}, changes={update:?})"
1967 );
1968 return;
1969 };
1970
1971 let new_since = collection.shared.lock_read_capabilities(|caps| {
1972 let read_frontier = caps.frontier();
1975 for (time, diff) in update.iter() {
1976 let count = caps.count_for(time) + diff;
1977 assert!(
1978 count >= 0,
1979 "invalid read capabilities update: negative capability \
1980 (id={id:?}, read_capabilities={caps:?}, update={update:?})",
1981 );
1982 assert!(
1983 count == 0 || read_frontier.less_equal(time),
1984 "invalid read capabilities update: frontier regression \
1985 (id={id:?}, read_capabilities={caps:?}, update={update:?})",
1986 );
1987 }
1988
1989 let changes = caps.update_iter(update.drain());
1992
1993 let changed = changes.count() > 0;
1994 changed.then(|| caps.frontier().to_owned())
1995 });
1996
1997 let Some(new_since) = new_since else {
1998 return; };
2000
2001 for read_hold in collection.compute_dependencies.values_mut() {
2003 read_hold
2004 .try_downgrade(new_since.clone())
2005 .expect("frontiers don't regress");
2006 }
2007 for read_hold in collection.storage_dependencies.values_mut() {
2008 read_hold
2009 .try_downgrade(new_since.clone())
2010 .expect("frontiers don't regress");
2011 }
2012
2013 self.send(ComputeCommand::AllowCompaction {
2015 id,
2016 frontier: new_since,
2017 });
2018 }
2019
2020 fn finish_peek(&mut self, uuid: Uuid, response: PeekResponse) {
2029 let Some(peek) = self.peeks.remove(&uuid) else {
2030 return;
2031 };
2032
2033 let _ = peek.peek_response_tx.send(response);
2035
2036 self.send(ComputeCommand::CancelPeek { uuid });
2039
2040 drop(peek.read_hold);
2041 }
2042
2043 fn handle_response(&mut self, (replica_id, epoch, response): ReplicaResponse) {
2046 if self
2048 .replicas
2049 .get(&replica_id)
2050 .filter(|replica| replica.epoch == epoch)
2051 .is_none()
2052 {
2053 return;
2054 }
2055
2056 match response {
2059 ComputeResponse::Frontiers(id, frontiers) => {
2060 self.handle_frontiers_response(id, frontiers, replica_id);
2061 }
2062 ComputeResponse::PeekResponse(uuid, peek_response, otel_ctx) => {
2063 self.handle_peek_response(uuid, peek_response, otel_ctx, replica_id);
2064 }
2065 ComputeResponse::CopyToResponse(id, response) => {
2066 self.handle_copy_to_response(id, response, replica_id);
2067 }
2068 ComputeResponse::SubscribeResponse(id, response) => {
2069 self.handle_subscribe_response(id, response, replica_id);
2070 }
2071 ComputeResponse::Status(response) => {
2072 self.handle_status_response(response, replica_id);
2073 }
2074 }
2075 }
2076
2077 fn handle_frontiers_response(
2080 &mut self,
2081 id: GlobalId,
2082 frontiers: FrontiersResponse,
2083 replica_id: ReplicaId,
2084 ) {
2085 if !self.collections.contains_key(&id) {
2086 soft_panic_or_log!(
2087 "frontiers update for an unknown collection \
2088 (id={id}, replica_id={replica_id}, frontiers={frontiers:?})"
2089 );
2090 return;
2091 }
2092 let Some(replica) = self.replicas.get_mut(&replica_id) else {
2093 soft_panic_or_log!(
2094 "frontiers update for an unknown replica \
2095 (replica_id={replica_id}, frontiers={frontiers:?})"
2096 );
2097 return;
2098 };
2099 let Some(replica_collection) = replica.collections.get_mut(&id) else {
2100 soft_panic_or_log!(
2101 "frontiers update for an unknown replica collection \
2102 (id={id}, replica_id={replica_id}, frontiers={frontiers:?})"
2103 );
2104 return;
2105 };
2106
2107 if let Some(new_frontier) = frontiers.input_frontier {
2108 replica_collection.update_input_frontier(new_frontier.clone());
2109 }
2110 if let Some(new_frontier) = frontiers.output_frontier {
2111 replica_collection.update_output_frontier(new_frontier.clone());
2112 }
2113 if let Some(new_frontier) = frontiers.write_frontier {
2114 replica_collection.update_write_frontier(new_frontier.clone());
2115 self.maybe_update_global_write_frontier(id, new_frontier);
2116 }
2117 }
2118
2119 #[mz_ore::instrument(level = "debug")]
2120 fn handle_peek_response(
2121 &mut self,
2122 uuid: Uuid,
2123 response: PeekResponse,
2124 otel_ctx: OpenTelemetryContext,
2125 replica_id: ReplicaId,
2126 ) {
2127 otel_ctx.attach_as_parent();
2128
2129 let Some(peek) = self.peeks.get(&uuid) else {
2132 return;
2133 };
2134
2135 let target_replica = peek.target_replica.unwrap_or(replica_id);
2137 if target_replica != replica_id {
2138 return;
2139 }
2140
2141 let duration = peek.requested_at.elapsed();
2142 self.metrics.observe_peek_response(&response, duration);
2143
2144 let notification = PeekNotification::new(&response, peek.offset, peek.limit);
2145 self.deliver_response(ComputeControllerResponse::PeekNotification(
2148 uuid,
2149 notification,
2150 otel_ctx,
2151 ));
2152
2153 self.finish_peek(uuid, response)
2154 }
2155
2156 fn handle_copy_to_response(
2157 &mut self,
2158 sink_id: GlobalId,
2159 response: CopyToResponse,
2160 replica_id: ReplicaId,
2161 ) {
2162 if !self.collections.contains_key(&sink_id) {
2163 soft_panic_or_log!(
2164 "received response for an unknown copy-to \
2165 (sink_id={sink_id}, replica_id={replica_id})",
2166 );
2167 return;
2168 }
2169 let Some(replica) = self.replicas.get_mut(&replica_id) else {
2170 soft_panic_or_log!("copy-to response for an unknown replica (replica_id={replica_id})");
2171 return;
2172 };
2173 let Some(replica_collection) = replica.collections.get_mut(&sink_id) else {
2174 soft_panic_or_log!(
2175 "copy-to response for an unknown replica collection \
2176 (sink_id={sink_id}, replica_id={replica_id})"
2177 );
2178 return;
2179 };
2180
2181 replica_collection.update_write_frontier(Antichain::new());
2185 replica_collection.update_input_frontier(Antichain::new());
2186 replica_collection.update_output_frontier(Antichain::new());
2187
2188 if !self.copy_tos.remove(&sink_id) {
2191 return;
2192 }
2193
2194 let result = match response {
2195 CopyToResponse::RowCount(count) => Ok(count),
2196 CopyToResponse::Error(error) => Err(anyhow::anyhow!(error)),
2197 CopyToResponse::Dropped => {
2202 tracing::error!(
2203 %sink_id, %replica_id,
2204 "received `Dropped` response for a tracked copy to",
2205 );
2206 return;
2207 }
2208 };
2209
2210 self.deliver_response(ComputeControllerResponse::CopyToResponse(sink_id, result));
2211 }
2212
2213 fn handle_subscribe_response(
2214 &mut self,
2215 subscribe_id: GlobalId,
2216 response: SubscribeResponse,
2217 replica_id: ReplicaId,
2218 ) {
2219 if !self.collections.contains_key(&subscribe_id) {
2220 soft_panic_or_log!(
2221 "received response for an unknown subscribe \
2222 (subscribe_id={subscribe_id}, replica_id={replica_id})",
2223 );
2224 return;
2225 }
2226 let Some(replica) = self.replicas.get_mut(&replica_id) else {
2227 soft_panic_or_log!(
2228 "subscribe response for an unknown replica (replica_id={replica_id})"
2229 );
2230 return;
2231 };
2232 let Some(replica_collection) = replica.collections.get_mut(&subscribe_id) else {
2233 soft_panic_or_log!(
2234 "subscribe response for an unknown replica collection \
2235 (subscribe_id={subscribe_id}, replica_id={replica_id})"
2236 );
2237 return;
2238 };
2239
2240 let write_frontier = match &response {
2244 SubscribeResponse::Batch(batch) => batch.upper.clone(),
2245 SubscribeResponse::DroppedAt(_) => Antichain::new(),
2246 };
2247
2248 replica_collection.update_write_frontier(write_frontier.clone());
2252 replica_collection.update_input_frontier(write_frontier.clone());
2253 replica_collection.update_output_frontier(write_frontier.clone());
2254
2255 let Some(mut subscribe) = self.subscribes.get(&subscribe_id).cloned() else {
2257 return;
2258 };
2259
2260 self.maybe_update_global_write_frontier(subscribe_id, write_frontier);
2266
2267 match response {
2268 SubscribeResponse::Batch(batch) => {
2269 let upper = batch.upper;
2270 let mut updates = batch.updates;
2271
2272 if PartialOrder::less_than(&subscribe.frontier, &upper) {
2275 let lower = std::mem::replace(&mut subscribe.frontier, upper.clone());
2276
2277 if upper.is_empty() {
2278 self.subscribes.remove(&subscribe_id);
2280 } else {
2281 self.subscribes.insert(subscribe_id, subscribe);
2283 }
2284
2285 if let Ok(updates) = updates.as_mut() {
2286 updates.retain_mut(|updates| {
2287 let offset = updates.times().partition_point(|t| {
2288 !lower.less_equal(t)
2291 });
2292 let (_, past_lower) = std::mem::take(updates).split_at(offset);
2293 *updates = past_lower;
2294 updates.len() > 0
2295 });
2296 }
2297 self.deliver_response(ComputeControllerResponse::SubscribeResponse(
2298 subscribe_id,
2299 SubscribeBatch {
2300 lower,
2301 upper,
2302 updates,
2303 },
2304 ));
2305 }
2306 }
2307 SubscribeResponse::DroppedAt(frontier) => {
2308 tracing::error!(
2313 %subscribe_id,
2314 %replica_id,
2315 frontier = ?frontier.elements(),
2316 "received `DroppedAt` response for a tracked subscribe",
2317 );
2318 self.subscribes.remove(&subscribe_id);
2319 }
2320 }
2321 }
2322
2323 fn handle_status_response(&self, response: StatusResponse, _replica_id: ReplicaId) {
2324 match response {
2325 StatusResponse::Placeholder => {}
2326 }
2327 }
2328
2329 fn dependency_write_frontiers<'b>(
2331 &'b self,
2332 collection: &'b CollectionState,
2333 ) -> impl Iterator<Item = Antichain<Timestamp>> + 'b {
2334 let compute_frontiers = collection.compute_dependency_ids().filter_map(|dep_id| {
2335 let collection = self.collections.get(&dep_id);
2336 collection.map(|c| c.write_frontier())
2337 });
2338 let storage_frontiers = collection.storage_dependency_ids().filter_map(|dep_id| {
2339 let frontiers = self.storage_collections.collection_frontiers(dep_id).ok();
2340 frontiers.map(|f| f.write_frontier)
2341 });
2342
2343 compute_frontiers.chain(storage_frontiers)
2344 }
2345
2346 fn transitive_storage_dependency_write_frontiers<'b>(
2348 &'b self,
2349 collection: &'b CollectionState,
2350 ) -> impl Iterator<Item = Antichain<Timestamp>> + 'b {
2351 let mut storage_ids: BTreeSet<_> = collection.storage_dependency_ids().collect();
2352 let mut todo: Vec<_> = collection.compute_dependency_ids().collect();
2353 let mut done = BTreeSet::new();
2354
2355 while let Some(id) = todo.pop() {
2356 if done.contains(&id) {
2357 continue;
2358 }
2359 if let Some(dep) = self.collections.get(&id) {
2360 storage_ids.extend(dep.storage_dependency_ids());
2361 todo.extend(dep.compute_dependency_ids())
2362 }
2363 done.insert(id);
2364 }
2365
2366 let storage_frontiers = storage_ids.into_iter().filter_map(|id| {
2367 let frontiers = self.storage_collections.collection_frontiers(id).ok();
2368 frontiers.map(|f| f.write_frontier)
2369 });
2370
2371 storage_frontiers
2372 }
2373
2374 fn downgrade_warmup_capabilities(&mut self) {
2387 let mut new_capabilities = BTreeMap::new();
2388 for (id, collection) in &self.collections {
2389 if collection.read_policy.is_none()
2393 && collection.shared.lock_write_frontier(|f| f.is_empty())
2394 {
2395 new_capabilities.insert(*id, Antichain::new());
2396 continue;
2397 }
2398
2399 let mut new_capability = Antichain::new();
2400 for frontier in self.dependency_write_frontiers(collection) {
2401 for time in frontier {
2402 new_capability.insert(time.step_back().unwrap_or(time));
2403 }
2404 }
2405
2406 new_capabilities.insert(*id, new_capability);
2407 }
2408
2409 for (id, new_capability) in new_capabilities {
2410 let collection = self.expect_collection_mut(id);
2411 let _ = collection.warmup_read_hold.try_downgrade(new_capability);
2412 }
2413 }
2414
2415 fn forward_implied_capabilities(&mut self) {
2443 if !ENABLE_PAUSED_CLUSTER_READHOLD_DOWNGRADE.get(&self.dyncfg) {
2444 return;
2445 }
2446 if !self.replicas.is_empty() {
2447 return;
2448 }
2449
2450 let mut new_capabilities = BTreeMap::new();
2451 for (id, collection) in &self.collections {
2452 let Some(read_policy) = &collection.read_policy else {
2453 continue;
2455 };
2456
2457 let mut dep_frontier = Antichain::new();
2461 for frontier in self.transitive_storage_dependency_write_frontiers(collection) {
2462 dep_frontier.extend(frontier);
2463 }
2464
2465 let new_capability = read_policy.frontier(dep_frontier.borrow());
2466 if PartialOrder::less_than(collection.implied_read_hold.since(), &new_capability) {
2467 new_capabilities.insert(*id, new_capability);
2468 }
2469 }
2470
2471 for (id, new_capability) in new_capabilities {
2472 let collection = self.expect_collection_mut(id);
2473 let _ = collection.implied_read_hold.try_downgrade(new_capability);
2474 }
2475 }
2476
2477 pub(super) fn acquire_read_hold(&self, id: GlobalId) -> Result<ReadHold, CollectionMissing> {
2482 let collection = self.collection(id)?;
2488 let since = collection.shared.lock_read_capabilities(|caps| {
2489 let since = caps.frontier().to_owned();
2490 caps.update_iter(since.iter().map(|t| (t.clone(), 1)));
2491 since
2492 });
2493 let hold = ReadHold::new(id, since, Arc::clone(&self.read_hold_tx));
2494 Ok(hold)
2495 }
2496
2497 #[mz_ore::instrument(level = "debug")]
2503 pub fn maintain(&mut self) {
2504 self.rehydrate_failed_replicas();
2505 self.downgrade_warmup_capabilities();
2506 self.forward_implied_capabilities();
2507 self.schedule_collections();
2508 self.cleanup_collections();
2509 self.update_frontier_introspection();
2510 self.refresh_state_metrics();
2511 self.refresh_wallclock_lag();
2512 }
2513}
2514
2515#[derive(Debug)]
2520struct CollectionState {
2521 target_replica: Option<ReplicaId>,
2523 log_collection: bool,
2527 dropped: bool,
2533 scheduled: bool,
2536
2537 read_only: bool,
2541
2542 shared: SharedCollectionState,
2544
2545 implied_read_hold: ReadHold,
2552 warmup_read_hold: ReadHold,
2560 read_policy: Option<ReadPolicy>,
2566
2567 storage_dependencies: BTreeMap<GlobalId, ReadHold>,
2570 compute_dependencies: BTreeMap<GlobalId, ReadHold>,
2573
2574 introspection: CollectionIntrospection,
2576
2577 wallclock_lag_histogram_stash: Option<
2584 BTreeMap<
2585 (
2586 WallclockLagHistogramPeriod,
2587 WallclockLag,
2588 BTreeMap<&'static str, String>,
2589 ),
2590 Diff,
2591 >,
2592 >,
2593}
2594
2595impl CollectionState {
2596 fn new(
2598 collection_id: GlobalId,
2599 as_of: Antichain<Timestamp>,
2600 shared: SharedCollectionState,
2601 storage_dependencies: BTreeMap<GlobalId, ReadHold>,
2602 compute_dependencies: BTreeMap<GlobalId, ReadHold>,
2603 read_hold_tx: read_holds::ChangeTx,
2604 introspection: CollectionIntrospection,
2605 ) -> Self {
2606 let since = as_of.clone();
2608 let upper = as_of;
2610
2611 assert!(shared.lock_read_capabilities(|c| c.frontier() == since.borrow()));
2613 assert!(shared.lock_write_frontier(|f| f == &upper));
2614
2615 let implied_read_hold =
2619 ReadHold::new(collection_id, since.clone(), Arc::clone(&read_hold_tx));
2620 let warmup_read_hold = ReadHold::new(collection_id, since.clone(), read_hold_tx);
2621
2622 let updates = warmup_read_hold.since().iter().map(|t| (t.clone(), 1));
2623 shared.lock_read_capabilities(|c| {
2624 c.update_iter(updates);
2625 });
2626
2627 let wallclock_lag_histogram_stash = match collection_id.is_transient() {
2631 true => None,
2632 false => Some(Default::default()),
2633 };
2634
2635 Self {
2636 target_replica: None,
2637 log_collection: false,
2638 dropped: false,
2639 scheduled: false,
2640 read_only: true,
2641 shared,
2642 implied_read_hold,
2643 warmup_read_hold,
2644 read_policy: Some(ReadPolicy::ValidFrom(since)),
2645 storage_dependencies,
2646 compute_dependencies,
2647 introspection,
2648 wallclock_lag_histogram_stash,
2649 }
2650 }
2651
2652 fn new_log_collection(
2654 id: GlobalId,
2655 shared: SharedCollectionState,
2656 read_hold_tx: read_holds::ChangeTx,
2657 introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
2658 ) -> Self {
2659 let since = Antichain::from_elem(Timestamp::MIN);
2660 let introspection = CollectionIntrospection::new(
2661 id,
2662 introspection_tx,
2663 since.clone(),
2664 false,
2665 None,
2666 None,
2667 Vec::new(),
2668 );
2669 let mut state = Self::new(
2670 id,
2671 since,
2672 shared,
2673 Default::default(),
2674 Default::default(),
2675 read_hold_tx,
2676 introspection,
2677 );
2678 state.log_collection = true;
2679 state.scheduled = true;
2681 state
2682 }
2683
2684 fn read_frontier(&self) -> Antichain<Timestamp> {
2686 self.shared
2687 .lock_read_capabilities(|c| c.frontier().to_owned())
2688 }
2689
2690 fn write_frontier(&self) -> Antichain<Timestamp> {
2692 self.shared.lock_write_frontier(|f| f.clone())
2693 }
2694
2695 fn storage_dependency_ids(&self) -> impl Iterator<Item = GlobalId> + '_ {
2696 self.storage_dependencies.keys().copied()
2697 }
2698
2699 fn compute_dependency_ids(&self) -> impl Iterator<Item = GlobalId> + '_ {
2700 self.compute_dependencies.keys().copied()
2701 }
2702}
2703
2704#[derive(Clone, Debug)]
2715pub(super) struct SharedCollectionState {
2716 read_capabilities: Arc<Mutex<MutableAntichain<Timestamp>>>,
2729 write_frontier: Arc<Mutex<Antichain<Timestamp>>>,
2731}
2732
2733impl SharedCollectionState {
2734 pub fn new(as_of: Antichain<Timestamp>) -> Self {
2735 let since = as_of.clone();
2737 let upper = as_of;
2739
2740 let mut read_capabilities = MutableAntichain::new();
2744 read_capabilities.update_iter(since.iter().map(|time| (time.clone(), 1)));
2745
2746 Self {
2747 read_capabilities: Arc::new(Mutex::new(read_capabilities)),
2748 write_frontier: Arc::new(Mutex::new(upper)),
2749 }
2750 }
2751
2752 pub fn lock_read_capabilities<F, R>(&self, f: F) -> R
2753 where
2754 F: FnOnce(&mut MutableAntichain<Timestamp>) -> R,
2755 {
2756 let mut caps = self.read_capabilities.lock().expect("poisoned");
2757 f(&mut *caps)
2758 }
2759
2760 pub fn lock_write_frontier<F, R>(&self, f: F) -> R
2761 where
2762 F: FnOnce(&mut Antichain<Timestamp>) -> R,
2763 {
2764 let mut frontier = self.write_frontier.lock().expect("poisoned");
2765 f(&mut *frontier)
2766 }
2767}
2768
2769#[derive(Debug)]
2772struct CollectionIntrospection {
2773 collection_id: GlobalId,
2775 introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
2777 frontiers: Option<FrontiersIntrospectionState>,
2782 refresh: Option<RefreshIntrospectionState>,
2786 dependency_ids: Vec<GlobalId>,
2788}
2789
2790impl CollectionIntrospection {
2791 fn new(
2792 collection_id: GlobalId,
2793 introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
2794 as_of: Antichain<Timestamp>,
2795 storage_sink: bool,
2796 initial_as_of: Option<Antichain<Timestamp>>,
2797 refresh_schedule: Option<RefreshSchedule>,
2798 dependency_ids: Vec<GlobalId>,
2799 ) -> Self {
2800 let refresh =
2801 match (refresh_schedule, initial_as_of) {
2802 (Some(refresh_schedule), Some(initial_as_of)) => Some(
2803 RefreshIntrospectionState::new(refresh_schedule, initial_as_of, &as_of),
2804 ),
2805 (refresh_schedule, _) => {
2806 soft_assert_or_log!(
2809 refresh_schedule.is_none(),
2810 "`refresh_schedule` without an `initial_as_of`: {collection_id}"
2811 );
2812 None
2813 }
2814 };
2815 let frontiers = (!storage_sink).then(|| FrontiersIntrospectionState::new(as_of));
2816
2817 let self_ = Self {
2818 collection_id,
2819 introspection_tx,
2820 frontiers,
2821 refresh,
2822 dependency_ids,
2823 };
2824
2825 self_.report_initial_state();
2826 self_
2827 }
2828
2829 fn report_initial_state(&self) {
2831 if let Some(frontiers) = &self.frontiers {
2832 let row = frontiers.row_for_collection(self.collection_id);
2833 let updates = vec![(row, Diff::ONE)];
2834 self.send(IntrospectionType::Frontiers, updates);
2835 }
2836
2837 if let Some(refresh) = &self.refresh {
2838 let row = refresh.row_for_collection(self.collection_id);
2839 let updates = vec![(row, Diff::ONE)];
2840 self.send(IntrospectionType::ComputeMaterializedViewRefreshes, updates);
2841 }
2842
2843 if !self.dependency_ids.is_empty() {
2844 let updates = self.dependency_rows(Diff::ONE);
2845 self.send(IntrospectionType::ComputeDependencies, updates);
2846 }
2847 }
2848
2849 fn dependency_rows(&self, diff: Diff) -> Vec<(Row, Diff)> {
2851 self.dependency_ids
2852 .iter()
2853 .map(|dependency_id| {
2854 let row = Row::pack_slice(&[
2855 Datum::String(&self.collection_id.to_string()),
2856 Datum::String(&dependency_id.to_string()),
2857 ]);
2858 (row, diff)
2859 })
2860 .collect()
2861 }
2862
2863 fn observe_frontiers(
2866 &mut self,
2867 read_frontier: &Antichain<Timestamp>,
2868 write_frontier: &Antichain<Timestamp>,
2869 ) {
2870 self.update_frontier_introspection(read_frontier, write_frontier);
2871 self.update_refresh_introspection(write_frontier);
2872 }
2873
2874 fn update_frontier_introspection(
2875 &mut self,
2876 read_frontier: &Antichain<Timestamp>,
2877 write_frontier: &Antichain<Timestamp>,
2878 ) {
2879 let Some(frontiers) = &mut self.frontiers else {
2880 return;
2881 };
2882
2883 if &frontiers.read_frontier == read_frontier && &frontiers.write_frontier == write_frontier
2884 {
2885 return; };
2887
2888 let retraction = frontiers.row_for_collection(self.collection_id);
2889 frontiers.update(read_frontier, write_frontier);
2890 let insertion = frontiers.row_for_collection(self.collection_id);
2891 let updates = vec![(retraction, Diff::MINUS_ONE), (insertion, Diff::ONE)];
2892 self.send(IntrospectionType::Frontiers, updates);
2893 }
2894
2895 fn update_refresh_introspection(&mut self, write_frontier: &Antichain<Timestamp>) {
2896 let Some(refresh) = &mut self.refresh else {
2897 return;
2898 };
2899
2900 let retraction = refresh.row_for_collection(self.collection_id);
2901 refresh.frontier_update(write_frontier);
2902 let insertion = refresh.row_for_collection(self.collection_id);
2903
2904 if retraction == insertion {
2905 return; }
2907
2908 let updates = vec![(retraction, Diff::MINUS_ONE), (insertion, Diff::ONE)];
2909 self.send(IntrospectionType::ComputeMaterializedViewRefreshes, updates);
2910 }
2911
2912 fn send(&self, introspection_type: IntrospectionType, updates: Vec<(Row, Diff)>) {
2913 let _ = self.introspection_tx.send((introspection_type, updates));
2916 }
2917}
2918
2919impl Drop for CollectionIntrospection {
2920 fn drop(&mut self) {
2921 if let Some(frontiers) = &self.frontiers {
2923 let row = frontiers.row_for_collection(self.collection_id);
2924 let updates = vec![(row, Diff::MINUS_ONE)];
2925 self.send(IntrospectionType::Frontiers, updates);
2926 }
2927
2928 if let Some(refresh) = &self.refresh {
2930 let retraction = refresh.row_for_collection(self.collection_id);
2931 let updates = vec![(retraction, Diff::MINUS_ONE)];
2932 self.send(IntrospectionType::ComputeMaterializedViewRefreshes, updates);
2933 }
2934
2935 if !self.dependency_ids.is_empty() {
2937 let updates = self.dependency_rows(Diff::MINUS_ONE);
2938 self.send(IntrospectionType::ComputeDependencies, updates);
2939 }
2940 }
2941}
2942
2943#[derive(Debug)]
2944struct FrontiersIntrospectionState {
2945 read_frontier: Antichain<Timestamp>,
2946 write_frontier: Antichain<Timestamp>,
2947}
2948
2949impl FrontiersIntrospectionState {
2950 fn new(as_of: Antichain<Timestamp>) -> Self {
2951 Self {
2952 read_frontier: as_of.clone(),
2953 write_frontier: as_of,
2954 }
2955 }
2956
2957 fn row_for_collection(&self, collection_id: GlobalId) -> Row {
2959 let read_frontier = self
2960 .read_frontier
2961 .as_option()
2962 .map_or(Datum::Null, |ts| ts.clone().into());
2963 let write_frontier = self
2964 .write_frontier
2965 .as_option()
2966 .map_or(Datum::Null, |ts| ts.clone().into());
2967 Row::pack_slice(&[
2968 Datum::String(&collection_id.to_string()),
2969 read_frontier,
2970 write_frontier,
2971 ])
2972 }
2973
2974 fn update(
2976 &mut self,
2977 read_frontier: &Antichain<Timestamp>,
2978 write_frontier: &Antichain<Timestamp>,
2979 ) {
2980 if read_frontier != &self.read_frontier {
2981 self.read_frontier.clone_from(read_frontier);
2982 }
2983 if write_frontier != &self.write_frontier {
2984 self.write_frontier.clone_from(write_frontier);
2985 }
2986 }
2987}
2988
2989#[derive(Debug)]
2992struct RefreshIntrospectionState {
2993 refresh_schedule: RefreshSchedule,
2995 initial_as_of: Antichain<Timestamp>,
2996 next_refresh: Datum<'static>, last_completed_refresh: Datum<'static>, }
3000
3001impl RefreshIntrospectionState {
3002 fn row_for_collection(&self, collection_id: GlobalId) -> Row {
3004 Row::pack_slice(&[
3005 Datum::String(&collection_id.to_string()),
3006 self.last_completed_refresh,
3007 self.next_refresh,
3008 ])
3009 }
3010}
3011
3012impl RefreshIntrospectionState {
3013 fn new(
3016 refresh_schedule: RefreshSchedule,
3017 initial_as_of: Antichain<Timestamp>,
3018 upper: &Antichain<Timestamp>,
3019 ) -> Self {
3020 let mut self_ = Self {
3021 refresh_schedule: refresh_schedule.clone(),
3022 initial_as_of: initial_as_of.clone(),
3023 next_refresh: Datum::Null,
3024 last_completed_refresh: Datum::Null,
3025 };
3026 self_.frontier_update(upper);
3027 self_
3028 }
3029
3030 fn frontier_update(&mut self, write_frontier: &Antichain<Timestamp>) {
3033 if write_frontier.is_empty() {
3034 self.last_completed_refresh =
3035 if let Some(last_refresh) = self.refresh_schedule.last_refresh() {
3036 last_refresh.into()
3037 } else {
3038 Timestamp::MAX.into()
3041 };
3042 self.next_refresh = Datum::Null;
3043 } else {
3044 if PartialOrder::less_equal(write_frontier, &self.initial_as_of) {
3045 self.last_completed_refresh = Datum::Null;
3047 let initial_as_of = self.initial_as_of.as_option().expect(
3048 "initial_as_of can't be [], because then there would be no refreshes at all",
3049 );
3050 let first_refresh = self
3051 .refresh_schedule
3052 .round_up_timestamp(*initial_as_of)
3053 .expect("sequencing makes sure that REFRESH MVs always have a first refresh");
3054 soft_assert_or_log!(
3055 first_refresh == *initial_as_of,
3056 "initial_as_of should be set to the first refresh"
3057 );
3058 self.next_refresh = first_refresh.into();
3059 } else {
3060 let write_frontier = write_frontier.as_option().expect("checked above");
3062 self.last_completed_refresh = self
3063 .refresh_schedule
3064 .round_down_timestamp_m1(*write_frontier)
3065 .map_or_else(
3066 || {
3067 soft_panic_or_log!(
3068 "rounding down should have returned the first refresh or later"
3069 );
3070 Datum::Null
3071 },
3072 |last_completed_refresh| last_completed_refresh.into(),
3073 );
3074 self.next_refresh = write_frontier.clone().into();
3075 }
3076 }
3077 }
3078}
3079
3080#[derive(Debug)]
3082struct PendingPeek {
3083 target_replica: Option<ReplicaId>,
3087 otel_ctx: OpenTelemetryContext,
3089 requested_at: Instant,
3093 read_hold: ReadHold,
3095 peek_response_tx: oneshot::Sender<PeekResponse>,
3097 limit: Option<usize>,
3099 offset: usize,
3101}
3102
3103#[derive(Debug, Clone)]
3104struct ActiveSubscribe {
3105 frontier: Antichain<Timestamp>,
3107}
3108
3109impl Default for ActiveSubscribe {
3110 fn default() -> Self {
3111 Self {
3112 frontier: Antichain::from_elem(Timestamp::MIN),
3113 }
3114 }
3115}
3116
3117#[derive(Debug)]
3119struct ReplicaState {
3120 id: ReplicaId,
3122 client: ReplicaClient,
3124 config: ReplicaConfig,
3126 metrics: ReplicaMetrics,
3128 introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
3130 collections: BTreeMap<GlobalId, ReplicaCollectionState>,
3132 epoch: u64,
3134}
3135
3136impl ReplicaState {
3137 fn new(
3138 id: ReplicaId,
3139 client: ReplicaClient,
3140 config: ReplicaConfig,
3141 metrics: ReplicaMetrics,
3142 introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
3143 epoch: u64,
3144 ) -> Self {
3145 Self {
3146 id,
3147 client,
3148 config,
3149 metrics,
3150 introspection_tx,
3151 epoch,
3152 collections: Default::default(),
3153 }
3154 }
3155
3156 fn add_collection(
3162 &mut self,
3163 id: GlobalId,
3164 as_of: Antichain<Timestamp>,
3165 input_read_holds: Vec<ReadHold>,
3166 ) {
3167 let metrics = self.metrics.for_collection(id);
3168 let introspection = ReplicaCollectionIntrospection::new(
3169 self.id,
3170 id,
3171 self.introspection_tx.clone(),
3172 as_of.clone(),
3173 );
3174 let mut state =
3175 ReplicaCollectionState::new(metrics, as_of, introspection, input_read_holds);
3176
3177 if id.is_transient() {
3181 state.wallclock_lag_max = None;
3182 }
3183
3184 if let Some(previous) = self.collections.insert(id, state) {
3185 panic!("attempt to add a collection with existing ID {id} (previous={previous:?}");
3186 }
3187 }
3188
3189 fn remove_collection(&mut self, id: GlobalId) -> Option<ReplicaCollectionState> {
3191 self.collections.remove(&id)
3192 }
3193
3194 fn expect_collection(&self, id: GlobalId) -> &ReplicaCollectionState {
3201 self.collections
3202 .get(&id)
3203 .expect("hosting replica must have per-replica collection state")
3204 }
3205
3206 fn collection_frontiers_empty(&self, id: GlobalId) -> bool {
3208 self.collections.get(&id).map_or(true, |c| {
3209 c.write_frontier.is_empty()
3210 && c.input_frontier.is_empty()
3211 && c.output_frontier.is_empty()
3212 })
3213 }
3214
3215 #[mz_ore::instrument(level = "debug")]
3219 pub fn dump(&self) -> Result<serde_json::Value, anyhow::Error> {
3220 let Self {
3227 id,
3228 client: _,
3229 config: _,
3230 metrics: _,
3231 introspection_tx: _,
3232 epoch,
3233 collections,
3234 } = self;
3235
3236 let collections: BTreeMap<_, _> = collections
3237 .iter()
3238 .map(|(id, collection)| (id.to_string(), format!("{collection:?}")))
3239 .collect();
3240
3241 Ok(serde_json::json!({
3242 "id": id.to_string(),
3243 "collections": collections,
3244 "epoch": epoch,
3245 }))
3246 }
3247}
3248
3249#[derive(Debug)]
3250struct ReplicaCollectionState {
3251 write_frontier: Antichain<Timestamp>,
3255 input_frontier: Antichain<Timestamp>,
3259 output_frontier: Antichain<Timestamp>,
3263
3264 metrics: Option<ReplicaCollectionMetrics>,
3268 as_of: Antichain<Timestamp>,
3270 introspection: ReplicaCollectionIntrospection,
3272 input_read_holds: Vec<ReadHold>,
3278
3279 wallclock_lag_max: Option<WallclockLag>,
3283}
3284
3285impl ReplicaCollectionState {
3286 fn new(
3287 metrics: Option<ReplicaCollectionMetrics>,
3288 as_of: Antichain<Timestamp>,
3289 introspection: ReplicaCollectionIntrospection,
3290 input_read_holds: Vec<ReadHold>,
3291 ) -> Self {
3292 Self {
3293 write_frontier: as_of.clone(),
3294 input_frontier: as_of.clone(),
3295 output_frontier: as_of.clone(),
3296 metrics,
3297 as_of,
3298 introspection,
3299 input_read_holds,
3300 wallclock_lag_max: Some(WallclockLag::MIN),
3301 }
3302 }
3303
3304 fn hydrated(&self) -> bool {
3306 self.as_of.is_empty() || PartialOrder::less_than(&self.as_of, &self.output_frontier)
3322 }
3323
3324 fn update_write_frontier(&mut self, new_frontier: Antichain<Timestamp>) {
3326 if PartialOrder::less_than(&new_frontier, &self.write_frontier) {
3327 soft_panic_or_log!(
3328 "replica collection write frontier regression (old={:?}, new={new_frontier:?})",
3329 self.write_frontier,
3330 );
3331 return;
3332 } else if new_frontier == self.write_frontier {
3333 return;
3334 }
3335
3336 self.write_frontier = new_frontier;
3337 }
3338
3339 fn update_input_frontier(&mut self, new_frontier: Antichain<Timestamp>) {
3341 if PartialOrder::less_than(&new_frontier, &self.input_frontier) {
3342 soft_panic_or_log!(
3343 "replica collection input frontier regression (old={:?}, new={new_frontier:?})",
3344 self.input_frontier,
3345 );
3346 return;
3347 } else if new_frontier == self.input_frontier {
3348 return;
3349 }
3350
3351 self.input_frontier = new_frontier;
3352
3353 for read_hold in &mut self.input_read_holds {
3355 let result = read_hold.try_downgrade(self.input_frontier.clone());
3356 soft_assert_or_log!(
3357 result.is_ok(),
3358 "read hold downgrade failed (read_hold={read_hold:?}, new_since={:?})",
3359 self.input_frontier,
3360 );
3361 }
3362 }
3363
3364 fn update_output_frontier(&mut self, new_frontier: Antichain<Timestamp>) {
3366 if PartialOrder::less_than(&new_frontier, &self.output_frontier) {
3367 soft_panic_or_log!(
3368 "replica collection output frontier regression (old={:?}, new={new_frontier:?})",
3369 self.output_frontier,
3370 );
3371 return;
3372 } else if new_frontier == self.output_frontier {
3373 return;
3374 }
3375
3376 self.output_frontier = new_frontier;
3377 }
3378}
3379
3380#[derive(Debug)]
3383struct ReplicaCollectionIntrospection {
3384 replica_id: ReplicaId,
3386 collection_id: GlobalId,
3388 write_frontier: Antichain<Timestamp>,
3390 introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
3392}
3393
3394impl ReplicaCollectionIntrospection {
3395 fn new(
3397 replica_id: ReplicaId,
3398 collection_id: GlobalId,
3399 introspection_tx: mpsc::UnboundedSender<IntrospectionUpdates>,
3400 as_of: Antichain<Timestamp>,
3401 ) -> Self {
3402 let self_ = Self {
3403 replica_id,
3404 collection_id,
3405 write_frontier: as_of,
3406 introspection_tx,
3407 };
3408
3409 self_.report_initial_state();
3410 self_
3411 }
3412
3413 fn report_initial_state(&self) {
3415 let row = self.write_frontier_row();
3416 let updates = vec![(row, Diff::ONE)];
3417 self.send(IntrospectionType::ReplicaFrontiers, updates);
3418 }
3419
3420 fn observe_frontier(&mut self, write_frontier: &Antichain<Timestamp>) {
3422 if self.write_frontier == *write_frontier {
3423 return; }
3425
3426 let retraction = self.write_frontier_row();
3427 self.write_frontier.clone_from(write_frontier);
3428 let insertion = self.write_frontier_row();
3429
3430 let updates = vec![(retraction, Diff::MINUS_ONE), (insertion, Diff::ONE)];
3431 self.send(IntrospectionType::ReplicaFrontiers, updates);
3432 }
3433
3434 fn write_frontier_row(&self) -> Row {
3436 let write_frontier = self
3437 .write_frontier
3438 .as_option()
3439 .map_or(Datum::Null, |ts| ts.clone().into());
3440 Row::pack_slice(&[
3441 Datum::String(&self.collection_id.to_string()),
3442 Datum::String(&self.replica_id.to_string()),
3443 write_frontier,
3444 ])
3445 }
3446
3447 fn send(&self, introspection_type: IntrospectionType, updates: Vec<(Row, Diff)>) {
3448 let _ = self.introspection_tx.send((introspection_type, updates));
3451 }
3452}
3453
3454impl Drop for ReplicaCollectionIntrospection {
3455 fn drop(&mut self) {
3456 let row = self.write_frontier_row();
3458 let updates = vec![(row, Diff::MINUS_ONE)];
3459 self.send(IntrospectionType::ReplicaFrontiers, updates);
3460 }
3461}
3462
3463fn classify_collection_readiness<'a, I>(
3469 replicas: I,
3470 target_replica_ids: &BTreeSet<ReplicaId>,
3471 reference_replica_ids: &BTreeSet<ReplicaId>,
3472 allowed_lag: Option<Timestamp>,
3473) -> CollectionReadiness
3474where
3475 I: Iterator<Item = (ReplicaId, &'a ReplicaCollectionState)> + Clone,
3476{
3477 let lag = allowed_lag.map(|allowed_lag| {
3478 let mut reference = Antichain::from_elem(Timestamp::MIN);
3479 for (_, state) in replicas
3480 .clone()
3481 .filter(|(id, _)| reference_replica_ids.contains(id))
3482 {
3483 reference.join_assign(&state.output_frontier);
3484 }
3485 (reference, allowed_lag)
3486 });
3487
3488 let mut result = CollectionReadiness::Unhydrated;
3489 for (_, state) in replicas.filter(|(id, _)| target_replica_ids.contains(id)) {
3490 match CollectionReadiness::classify(
3491 state.hydrated(),
3492 &state.output_frontier,
3493 lag.as_ref().map(|(reference, lag)| (reference, *lag)),
3494 ) {
3495 CollectionReadiness::Ready => return CollectionReadiness::Ready,
3496 CollectionReadiness::Lagging { lag } => {
3497 let lag = match result {
3500 CollectionReadiness::Lagging {
3501 lag: Some(previous),
3502 } => Some(previous.min(lag.unwrap_or(u64::MAX))),
3503 _ => lag,
3504 };
3505 result = CollectionReadiness::Lagging { lag };
3506 }
3507 CollectionReadiness::Unhydrated => {}
3508 }
3509 }
3510 result
3511}
3512
3513#[cfg(test)]
3514mod tests {
3515 use std::collections::{BTreeMap, BTreeSet};
3516
3517 use mz_compute_types::dyncfgs::{ENABLE_COLUMN_PAGED_BATCHER, ENABLE_MZ_JOIN_CORE};
3518 use mz_dyncfg::{ConfigSet, ConfigUpdates, ConfigVal};
3519 use mz_persist_types::PersistLocation;
3520 use mz_repr::{GlobalId, Timestamp};
3521 use timely::progress::Antichain;
3522 use tokio::sync::mpsc;
3523
3524 use crate::protocol::command::{ComputeCommand, InstanceConfig};
3525
3526 use super::{
3527 CollectionReadiness, Instance, ReplicaCollectionIntrospection, ReplicaCollectionState,
3528 ReplicaId, classify_collection_readiness,
3529 };
3530
3531 fn ac(ts: u64) -> Antichain<Timestamp> {
3532 Antichain::from_elem(Timestamp::new(ts))
3533 }
3534
3535 fn state(id: ReplicaId, as_of: u64, write: u64, output: u64) -> ReplicaCollectionState {
3536 let (tx, _rx) = mpsc::unbounded_channel();
3537 let introspection =
3538 ReplicaCollectionIntrospection::new(id, GlobalId::User(1), tx, ac(as_of));
3539 let mut state = ReplicaCollectionState::new(None, ac(as_of), introspection, Vec::new());
3540 state.update_write_frontier(ac(write));
3541 state.update_output_frontier(ac(output));
3542 state
3543 }
3544
3545 fn classify(
3546 replicas: &[(ReplicaId, ReplicaCollectionState)],
3547 targets: &[ReplicaId],
3548 references: &[ReplicaId],
3549 lag: Option<Timestamp>,
3550 ) -> CollectionReadiness {
3551 classify_collection_readiness(
3552 replicas.iter().map(|(id, state)| (*id, state)),
3553 &targets.iter().copied().collect::<BTreeSet<_>>(),
3554 &references.iter().copied().collect::<BTreeSet<_>>(),
3555 lag,
3556 )
3557 }
3558
3559 const LAG: Option<Timestamp> = Some(Timestamp::new(60));
3560
3561 #[mz_ore::test]
3565 fn shared_write_does_not_hide_output_lag() {
3566 let reference = ReplicaId::User(1);
3567 let target = ReplicaId::User(2);
3568 let mut replicas = vec![
3569 (reference, state(reference, 100, 10_000, 10_000)),
3570 (target, state(target, 100, 10_000, 101)),
3573 ];
3574 assert_eq!(
3575 classify(&replicas, &[target], &[reference], LAG),
3576 CollectionReadiness::Lagging { lag: Some(9_899) },
3577 );
3578 replicas[1].1.update_output_frontier(ac(9_940));
3579 assert_eq!(
3580 classify(&replicas, &[target], &[reference], LAG),
3581 CollectionReadiness::Ready,
3582 );
3583 }
3584
3585 #[mz_ore::test]
3586 fn refresh_writes_ahead_of_outputs_is_ready() {
3587 let reference = ReplicaId::User(1);
3588 let target = ReplicaId::User(2);
3589 let replicas = vec![
3590 (reference, state(reference, 100, 20_000, 10_000)),
3591 (target, state(target, 100, 20_000, 10_000)),
3592 ];
3593 assert_eq!(
3594 classify(&replicas, &[target], &[reference], LAG),
3595 CollectionReadiness::Ready,
3596 );
3597 }
3598
3599 #[mz_ore::test]
3600 fn furthest_reference_and_bystander_exclusion() {
3601 let trailing = ReplicaId::User(1);
3602 let ahead = ReplicaId::User(2);
3603 let target = ReplicaId::User(3);
3604 let bystander = ReplicaId::User(4);
3605 let replicas = vec![
3606 (trailing, state(trailing, 100, 9_000, 9_000)),
3607 (ahead, state(ahead, 100, 10_000, 10_000)),
3608 (target, state(target, 100, 8_990, 8_990)),
3609 (bystander, state(bystander, 100, 1_000_000, 1_000_000)),
3610 ];
3611 assert_eq!(
3612 classify(&replicas, &[target], &[trailing], LAG),
3613 CollectionReadiness::Ready,
3614 );
3615 assert_eq!(
3616 classify(&replicas, &[target], &[trailing, ahead], LAG),
3617 CollectionReadiness::Lagging { lag: Some(1_010) },
3618 );
3619 }
3620
3621 #[mz_ore::test]
3622 fn one_ready_target_suffices() {
3623 let reference = ReplicaId::User(1);
3624 let unhydrated = ReplicaId::User(2);
3625 let lagging = ReplicaId::User(3);
3626 let ready = ReplicaId::User(4);
3627 let replicas = vec![
3628 (reference, state(reference, 100, 10_000, 10_000)),
3629 (unhydrated, state(unhydrated, 10_000, 10_000, 10_000)),
3633 (lagging, state(lagging, 100, 1_000, 1_000)),
3634 (ready, state(ready, 100, 9_990, 9_990)),
3635 ];
3636 assert_eq!(
3637 classify(&replicas, &[unhydrated, lagging, ready], &[reference], LAG),
3638 CollectionReadiness::Ready,
3639 );
3640 assert_eq!(
3641 classify(&replicas, &[unhydrated, lagging], &[reference], LAG),
3642 CollectionReadiness::Lagging { lag: Some(9_000) },
3643 );
3644 }
3645
3646 #[mz_ore::test]
3647 fn no_lag_gate_is_hydration_only() {
3648 let reference = ReplicaId::User(1);
3649 let target = ReplicaId::User(2);
3650 let replicas = vec![
3651 (reference, state(reference, 100, 10_000, 10_000)),
3652 (target, state(target, 100, 101, 101)),
3653 ];
3654 assert_eq!(
3655 classify(&replicas, &[target], &[reference], None),
3656 CollectionReadiness::Ready,
3657 );
3658 let unhydrated = vec![(target, state(target, 100, 100, 100))];
3659 assert_eq!(
3660 classify(&unhydrated, &[target], &[], None),
3661 CollectionReadiness::Unhydrated,
3662 );
3663 }
3664
3665 #[mz_ore::test]
3666 fn lag_reports_the_closest_hydrated_target() {
3667 let reference = ReplicaId::User(1);
3668 let far = ReplicaId::User(2);
3669 let close = ReplicaId::User(3);
3670 let mut replicas = vec![
3671 (reference, state(reference, 100, 10_000, 10_000)),
3672 (far, state(far, 100, 1_000, 1_000)),
3673 (close, state(close, 100, 9_000, 9_000)),
3674 ];
3675 for _ in 0..2 {
3676 assert_eq!(
3677 classify(&replicas, &[far, close], &[reference], LAG),
3678 CollectionReadiness::Lagging { lag: Some(1_000) },
3679 );
3680 replicas.reverse();
3681 }
3682 }
3683
3684 #[mz_ore::test]
3685 fn completed_reference_requires_completed_target() {
3686 let reference = ReplicaId::User(1);
3687 let target = ReplicaId::User(2);
3688 let mut replicas = vec![
3689 (reference, state(reference, 100, 1_000, 1_000)),
3690 (target, state(target, 100, 1_000, 1_000)),
3691 ];
3692 replicas[0].1.update_output_frontier(Antichain::new());
3693 assert_eq!(
3694 classify(&replicas, &[target], &[reference], Some(Timestamp::new(0))),
3695 CollectionReadiness::Lagging { lag: None },
3696 );
3697 replicas[1].1.update_output_frontier(Antichain::new());
3698 assert_eq!(
3699 classify(&replicas, &[target], &[reference], Some(Timestamp::new(0))),
3700 CollectionReadiness::Ready,
3701 );
3702 }
3703
3704 #[mz_ore::test]
3705 fn no_reference_means_nothing_to_regress() {
3706 let target = ReplicaId::User(1);
3707 let replicas = vec![(target, state(target, 100, 5_000, 5_000))];
3708 assert_eq!(
3709 classify(&replicas, &[target], &[], Some(Timestamp::new(0))),
3710 CollectionReadiness::Ready,
3711 );
3712 assert_eq!(
3713 classify(&replicas, &[target], &[target], Some(Timestamp::new(0))),
3714 CollectionReadiness::Ready,
3715 );
3716 }
3717
3718 #[mz_ore::test]
3719 fn no_targets_is_unhydrated() {
3720 let reference = ReplicaId::User(1);
3721 let replicas = vec![(reference, state(reference, 100, 10_000, 10_000))];
3722 assert_eq!(
3723 classify(&replicas, &[], &[reference], LAG),
3724 CollectionReadiness::Unhydrated,
3725 );
3726 }
3727
3728 fn create_instance_command() -> ComputeCommand {
3729 ComputeCommand::CreateInstance(Box::new(InstanceConfig {
3730 logging: Default::default(),
3731 expiration_offset: None,
3732 peek_stash_persist_location: PersistLocation::new_in_mem(),
3733 arrangement_dictionary_compression: false,
3734 initial_config: Default::default(),
3735 }))
3736 }
3737
3738 fn initial_config(cmd: &ComputeCommand) -> &ConfigUpdates {
3739 match cmd {
3740 ComputeCommand::CreateInstance(config) => &config.initial_config,
3741 other => panic!("expected CreateInstance, got {other:?}"),
3742 }
3743 }
3744
3745 #[mz_ore::test]
3750 fn create_instance_snapshots_instance_wide_dyncfg() {
3751 let dyncfg = ConfigSet::default()
3752 .add(&ENABLE_COLUMN_PAGED_BATCHER)
3753 .add(&ENABLE_MZ_JOIN_CORE);
3754 let mut updates = ConfigUpdates::default();
3755 updates.add(&ENABLE_COLUMN_PAGED_BATCHER, true);
3756 updates.add(&ENABLE_MZ_JOIN_CORE, false);
3757 updates.apply(&dyncfg);
3758
3759 let overrides = BTreeMap::new();
3761 let cmd = Instance::specialize_command_for_replica(
3762 create_instance_command(),
3763 ReplicaId::User(1),
3764 &overrides,
3765 &dyncfg,
3766 );
3767 let snapshot = initial_config(&cmd);
3768 assert_eq!(
3769 snapshot.updates.get(ENABLE_COLUMN_PAGED_BATCHER.name()),
3770 Some(&ConfigVal::Bool(true)),
3771 );
3772 assert_eq!(
3773 snapshot.updates.get(ENABLE_MZ_JOIN_CORE.name()),
3774 Some(&ConfigVal::Bool(false)),
3775 );
3776 }
3777
3778 #[mz_ore::test]
3781 fn create_instance_snapshot_applies_replica_override() {
3782 let dyncfg = ConfigSet::default().add(&ENABLE_COLUMN_PAGED_BATCHER);
3783 let mut updates = ConfigUpdates::default();
3784 updates.add(&ENABLE_COLUMN_PAGED_BATCHER, true);
3785 updates.apply(&dyncfg);
3786
3787 let replica = ReplicaId::User(1);
3788 let mut override_updates = ConfigUpdates::default();
3789 override_updates.add(&ENABLE_COLUMN_PAGED_BATCHER, false);
3790 let overrides = BTreeMap::from([(replica, override_updates)]);
3791
3792 let cmd = Instance::specialize_command_for_replica(
3793 create_instance_command(),
3794 replica,
3795 &overrides,
3796 &dyncfg,
3797 );
3798 assert_eq!(
3799 initial_config(&cmd)
3800 .updates
3801 .get(ENABLE_COLUMN_PAGED_BATCHER.name()),
3802 Some(&ConfigVal::Bool(false)),
3803 "replica override should win over the instance-wide value",
3804 );
3805 }
3806
3807 #[mz_ore::test]
3809 fn update_configuration_merges_replica_override() {
3810 let dyncfg = ConfigSet::default().add(&ENABLE_COLUMN_PAGED_BATCHER);
3811
3812 let replica = ReplicaId::User(1);
3813 let mut override_updates = ConfigUpdates::default();
3814 override_updates.add(&ENABLE_COLUMN_PAGED_BATCHER, true);
3815 let overrides = BTreeMap::from([(replica, override_updates)]);
3816
3817 let cmd = Instance::specialize_command_for_replica(
3818 ComputeCommand::UpdateConfiguration(Box::new(Default::default())),
3819 replica,
3820 &overrides,
3821 &dyncfg,
3822 );
3823 match cmd {
3824 ComputeCommand::UpdateConfiguration(params) => assert_eq!(
3825 params
3826 .dyncfg_updates
3827 .updates
3828 .get(ENABLE_COLUMN_PAGED_BATCHER.name()),
3829 Some(&ConfigVal::Bool(true)),
3830 ),
3831 other => panic!("expected UpdateConfiguration, got {other:?}"),
3832 }
3833 }
3834}