1use std::collections::{BTreeMap, BTreeSet};
14use std::pin::Pin;
15use std::sync::Arc;
16use std::time::{Duration, Instant};
17
18use fail::fail_point;
19use maplit::{btreemap, btreeset};
20use mz_adapter_types::compaction::SINCE_GRANULARITY;
21use mz_adapter_types::connection::ConnectionId;
22use mz_audit_log::VersionedEvent;
23use mz_catalog::SYSTEM_CONN_ID;
24use mz_catalog::memory::objects::{CatalogItem, DataSourceDesc, Sink};
25use mz_cluster_client::ReplicaId;
26use mz_controller::clusters::ReplicaLocation;
27use mz_controller_types::ClusterId;
28use mz_ore::instrument;
29use mz_ore::metrics::MetricsFutureExt;
30use mz_ore::now::to_datetime;
31use mz_ore::retry::Retry;
32use mz_ore::task;
33use mz_repr::adt::numeric::Numeric;
34use mz_repr::{CatalogItemId, GlobalId, Timestamp};
35use mz_sql::catalog::{CatalogClusterReplica, CatalogSchema};
36use mz_sql::names::ResolvedDatabaseSpecifier;
37use mz_sql::plan::ConnectionDetails;
38use mz_sql::session::metadata::SessionMetadata;
39use mz_sql::session::vars::{
40 self, DEFAULT_TIMESTAMP_INTERVAL, MAX_AWS_PRIVATELINK_CONNECTIONS, MAX_CLUSTERS,
41 MAX_CREDIT_CONSUMPTION_RATE, MAX_DATABASES, MAX_KAFKA_CONNECTIONS, MAX_MATERIALIZED_VIEWS,
42 MAX_MYSQL_CONNECTIONS, MAX_NETWORK_POLICIES, MAX_OBJECTS_PER_SCHEMA, MAX_POSTGRES_CONNECTIONS,
43 MAX_REPLICAS_PER_CLUSTER, MAX_ROLES, MAX_SCHEMAS_PER_DATABASE, MAX_SECRETS, MAX_SINKS,
44 MAX_SOURCES, MAX_SQL_SERVER_CONNECTIONS, MAX_TABLES, SystemVars, Var,
45};
46use mz_storage_client::controller::{CollectionDescription, DataSource, ExportDescription};
47use mz_storage_types::connections::inline::IntoInlineConnection;
48use mz_storage_types::read_policy::ReadPolicy;
49use mz_storage_types::sources::kafka::KAFKA_PROGRESS_DESC;
50use serde_json::json;
51use tracing::{Instrument, Level, event, info_span, warn};
52
53use crate::active_compute_sink::{ActiveComputeSink, ActiveComputeSinkRetireReason};
54use crate::catalog::{DropObjectInfo, Op, TransactionResult};
55use crate::coord::Coordinator;
56use crate::coord::appends::{BuiltinTableAppendCompletion, BuiltinTableAppendNotify};
57use crate::coord::catalog_implications::parsed_state_updates::ParsedStateUpdate;
58use crate::session::{Session, Transaction, TransactionOps};
59use crate::telemetry::{EventDetails, SegmentClientExt};
60use crate::util::ResultExt;
61use crate::{AdapterError, ExecuteContext, catalog, flags};
62
63impl Coordinator {
64 #[instrument(name = "coord::catalog_transact")]
66 pub(crate) async fn catalog_transact(
67 &mut self,
68 session: Option<&Session>,
69 ops: Vec<catalog::Op>,
70 ) -> Result<(), AdapterError> {
71 let start = Instant::now();
72 let result = self
73 .catalog_transact_with_context(session.map(|session| session.conn_id()), None, ops)
74 .await;
75 self.metrics
76 .catalog_transact_seconds
77 .with_label_values(&["catalog_transact"])
78 .observe(start.elapsed().as_secs_f64());
79 result
80 }
81
82 #[instrument(name = "coord::catalog_transact_with_side_effects")]
91 pub(crate) async fn catalog_transact_with_side_effects<F>(
92 &mut self,
93 mut ctx: Option<&mut ExecuteContext>,
94 ops: Vec<catalog::Op>,
95 side_effect: F,
96 ) -> Result<(), AdapterError>
97 where
98 F: for<'a> FnOnce(
99 &'a mut Coordinator,
100 Option<&'a mut ExecuteContext>,
101 ) -> Pin<Box<dyn Future<Output = ()> + 'a>>
102 + 'static,
103 {
104 let start = Instant::now();
105
106 let (table_updates, catalog_updates) = self
107 .catalog_transact_inner(ctx.as_ref().map(|ctx| ctx.session().conn_id()), ops)
108 .await?;
109
110 let apply_implications_res = self
113 .apply_catalog_implications(ctx.as_deref_mut(), catalog_updates)
114 .await;
115
116 apply_implications_res.expect("cannot fail to apply catalog update implications");
120
121 let consistency_start = Instant::now();
128
129 mz_ore::soft_assert_eq_no_log!(
132 self.check_consistency(),
133 Ok(()),
134 "coordinator inconsistency detected"
135 );
136
137 self.metrics
138 .catalog_transact_phase_seconds
139 .with_label_values(&["consistency_check"])
140 .observe(consistency_start.elapsed().as_secs_f64());
141
142 let side_effects_seconds = self
143 .metrics
144 .catalog_transact_phase_seconds
145 .with_label_values(&["side_effects"]);
146 let table_updates_wait = self
151 .metrics
152 .catalog_transact_phase_seconds
153 .with_label_values(&["table_updates_residual_wait"]);
154 let side_effects_fut = side_effect(self, ctx);
155
156 let ((), ()) = futures::future::join(
158 side_effects_fut
159 .wall_time()
160 .observe(side_effects_seconds)
161 .instrument(info_span!(
162 "coord::catalog_transact_with_side_effects::side_effects_fut"
163 )),
164 table_updates
165 .wall_time()
166 .observe(table_updates_wait)
167 .instrument(info_span!(
168 "coord::catalog_transact_with_side_effects::table_updates"
169 )),
170 )
171 .await;
172
173 self.metrics
174 .catalog_transact_seconds
175 .with_label_values(&["catalog_transact_with_side_effects"])
176 .observe(start.elapsed().as_secs_f64());
177
178 Ok(())
179 }
180
181 #[instrument(name = "coord::catalog_transact_with_context")]
189 pub(crate) async fn catalog_transact_with_context(
190 &mut self,
191 conn_id: Option<&ConnectionId>,
192 ctx: Option<&mut ExecuteContext>,
193 ops: Vec<catalog::Op>,
194 ) -> Result<(), AdapterError> {
195 let start = Instant::now();
196
197 let conn_id = conn_id.or_else(|| ctx.as_ref().map(|ctx| ctx.session().conn_id()));
198
199 let (table_updates, catalog_updates) = self.catalog_transact_inner(conn_id, ops).await?;
200
201 let table_updates_wait = self
202 .metrics
203 .catalog_transact_phase_seconds
204 .with_label_values(&["table_updates_wait"]);
205 let apply_catalog_implications_fut = self.apply_catalog_implications(ctx, catalog_updates);
206
207 let (combined_apply_res, ()) = futures::future::join(
209 apply_catalog_implications_fut.instrument(info_span!(
210 "coord::catalog_transact_with_context::side_effects_fut"
211 )),
212 table_updates
213 .wall_time()
214 .observe(table_updates_wait)
215 .instrument(info_span!(
216 "coord::catalog_transact_with_context::table_updates"
217 )),
218 )
219 .await;
220
221 combined_apply_res.expect("cannot fail to apply catalog implications");
225
226 let consistency_start = Instant::now();
229
230 mz_ore::soft_assert_eq_no_log!(
233 self.check_consistency(),
234 Ok(()),
235 "coordinator inconsistency detected"
236 );
237
238 self.metrics
239 .catalog_transact_phase_seconds
240 .with_label_values(&["consistency_check"])
241 .observe(consistency_start.elapsed().as_secs_f64());
242
243 self.metrics
244 .catalog_transact_seconds
245 .with_label_values(&["catalog_transact_with_context"])
246 .observe(start.elapsed().as_secs_f64());
247
248 Ok(())
249 }
250
251 #[instrument(name = "coord::catalog_transact_with_ddl_transaction")]
254 pub(crate) async fn catalog_transact_with_ddl_transaction<F>(
255 &mut self,
256 ctx: &mut ExecuteContext,
257 mut ops: Vec<catalog::Op>,
258 side_effect: F,
259 ) -> Result<(), AdapterError>
260 where
261 F: for<'a> FnOnce(
262 &'a mut Coordinator,
263 Option<&'a mut ExecuteContext>,
264 ) -> Pin<Box<dyn Future<Output = ()> + 'a>>
265 + Send
266 + Sync
267 + 'static,
268 {
269 let start = Instant::now();
270
271 let Some(Transaction {
272 ops:
273 TransactionOps::DDL {
274 ops: txn_ops,
275 revision: txn_revision,
276 state: txn_state,
277 snapshot: txn_snapshot,
278 side_effects: _,
279 },
280 ..
281 }) = ctx.session().transaction().inner()
282 else {
283 let result = self
284 .catalog_transact_with_side_effects(Some(ctx), ops, side_effect)
285 .await;
286 self.metrics
287 .catalog_transact_seconds
288 .with_label_values(&["catalog_transact_with_ddl_transaction"])
289 .observe(start.elapsed().as_secs_f64());
290 return result;
291 };
292
293 if self.catalog().transient_revision() != *txn_revision {
295 self.metrics
296 .catalog_transact_seconds
297 .with_label_values(&["catalog_transact_with_ddl_transaction"])
298 .observe(start.elapsed().as_secs_f64());
299 return Err(AdapterError::DDLTransactionRace);
300 }
301
302 let phase_seconds = self.metrics.catalog_transact_phase_seconds.clone();
307
308 let clone_start = Instant::now();
310 let txn_ops_clone = txn_ops.clone();
311 let txn_state_clone = txn_state.clone();
312 let prev_snapshot = txn_snapshot.clone();
316 phase_seconds
317 .with_label_values(&["ddl_txn_snapshot_clone"])
318 .observe(clone_start.elapsed().as_secs_f64());
319
320 let prep_start = Instant::now();
322 let mut combined_ops = txn_ops_clone;
323 combined_ops.extend(ops.iter().cloned());
324 let creates_scoped_object = ops.iter().any(|op| {
325 matches!(
326 op,
327 catalog::Op::CreateCluster { .. } | catalog::Op::CreateClusterReplica { .. }
328 )
329 });
330 if creates_scoped_object {
331 if let Some(scoped_op) = self.scoped_overrides_create_op(&combined_ops) {
335 ops.push(scoped_op.clone());
336 combined_ops.push(scoped_op);
337 }
338 }
339 let conn_id = ctx.session().conn_id().clone();
340 let validate_res = self.validate_resource_limits(&combined_ops, &conn_id);
341 phase_seconds
342 .with_label_values(&["ddl_txn_prep"])
343 .observe(prep_start.elapsed().as_secs_f64());
344 validate_res?;
345
346 let oracle_write_ts = self
348 .get_local_write_ts()
349 .wall_time()
350 .observe(phase_seconds.with_label_values(&["ddl_txn_write_ts"]))
351 .await
352 .timestamp;
353
354 let conn = self.active_conns.get(ctx.session().conn_id());
356
357 let (new_state, new_snapshot) = self
363 .catalog()
364 .transact_incremental_dry_run(
365 &txn_state_clone,
366 ops.clone(),
367 conn,
368 prev_snapshot,
369 oracle_write_ts,
370 )
371 .wall_time()
372 .observe(phase_seconds.with_label_values(&["ddl_txn_dry_run"]))
373 .await?;
374
375 let result = ctx
377 .session_mut()
378 .transaction_mut()
379 .add_ops(TransactionOps::DDL {
380 ops: combined_ops,
381 state: new_state,
382 side_effects: vec![Box::new(side_effect)],
383 revision: self.catalog().transient_revision(),
384 snapshot: Some(new_snapshot),
385 });
386
387 self.metrics
388 .catalog_transact_seconds
389 .with_label_values(&["catalog_transact_with_ddl_transaction"])
390 .observe(start.elapsed().as_secs_f64());
391
392 result
393 }
394
395 #[instrument(name = "coord::catalog_transact_inner")]
399 pub(crate) async fn catalog_transact_inner(
400 &mut self,
401 conn_id: Option<&ConnectionId>,
402 mut ops: Vec<catalog::Op>,
403 ) -> Result<(BuiltinTableAppendNotify, Vec<ParsedStateUpdate>), AdapterError> {
404 if self.controller.read_only() {
405 return Err(AdapterError::ReadOnly);
406 }
407
408 if let Some(scoped_op) = self.scoped_overrides_create_op(&ops) {
409 ops.push(scoped_op);
410 }
411
412 event!(Level::TRACE, ops = format!("{:?}", ops));
413
414 let phase_seconds = self.metrics.catalog_transact_phase_seconds.clone();
415 let phase_start = Instant::now();
416
417 let mut webhook_sources_to_restart = BTreeSet::new();
418 let mut clusters_to_drop = vec![];
419 let mut cluster_replicas_to_drop = vec![];
420 let mut clusters_to_create = vec![];
421 let mut cluster_replicas_to_create = vec![];
422 let mut update_metrics_config = false;
423 let mut update_tracing_config = false;
424 let mut update_controller_config = false;
425 let mut update_compute_config = false;
426 let mut update_storage_config = false;
427 let mut update_timestamp_oracle_config = false;
428 let mut update_metrics_retention = false;
429 let mut update_secrets_caching_config = false;
430 let mut update_cluster_scheduling_config = false;
431 let mut update_http_config = false;
432 let mut update_advance_timelines_interval = false;
433 let mut update_optimizer_e2e_latency_warning_threshold = false;
434 let mut reconcile_metric_sinks = false;
435
436 for op in &ops {
437 match op {
438 catalog::Op::DropObjects(drop_object_infos) => {
439 for drop_object_info in drop_object_infos {
440 match &drop_object_info {
441 catalog::DropObjectInfo::Item(_) => {
442 }
445 catalog::DropObjectInfo::Cluster(id) => {
446 clusters_to_drop.push(*id);
447 }
448 catalog::DropObjectInfo::ClusterReplica((
449 cluster_id,
450 replica_id,
451 _reason,
452 )) => {
453 cluster_replicas_to_drop.push((*cluster_id, *replica_id));
455 }
456 _ => (),
457 }
458 }
459 }
460 catalog::Op::ResetSystemConfiguration { name }
461 | catalog::Op::UpdateSystemConfiguration { name, .. } => {
462 update_metrics_config |= self
463 .catalog
464 .state()
465 .system_config()
466 .is_metrics_config_var(name);
467 update_tracing_config |= vars::is_tracing_var(name);
468 update_controller_config |= self
469 .catalog
470 .state()
471 .system_config()
472 .is_controller_config_var(name);
473 update_compute_config |= self
474 .catalog
475 .state()
476 .system_config()
477 .is_compute_config_var(name);
478 update_storage_config |= self
479 .catalog
480 .state()
481 .system_config()
482 .is_storage_config_var(name);
483 update_timestamp_oracle_config |= vars::is_timestamp_oracle_config_var(name);
484 update_metrics_retention |= name == vars::METRICS_RETENTION.name();
485 update_secrets_caching_config |= vars::is_secrets_caching_var(name);
486 update_cluster_scheduling_config |= vars::is_cluster_scheduling_var(name);
487 update_http_config |= vars::is_http_config_var(name);
488 update_advance_timelines_interval |= name == DEFAULT_TIMESTAMP_INTERVAL.name();
489 update_optimizer_e2e_latency_warning_threshold |=
490 name == vars::OPTIMIZER_E2E_LATENCY_WARNING_THRESHOLD.name();
491 reconcile_metric_sinks |= name == vars::DISABLED_METRIC_SINKS.name();
492 }
493 catalog::Op::ResetAllSystemConfiguration => {
494 update_tracing_config = true;
498 update_controller_config = true;
499 update_compute_config = true;
500 update_storage_config = true;
501 update_timestamp_oracle_config = true;
502 update_metrics_retention = true;
503 update_secrets_caching_config = true;
504 update_cluster_scheduling_config = true;
505 update_metrics_config = true;
506 update_http_config = true;
507 update_advance_timelines_interval = true;
508 update_optimizer_e2e_latency_warning_threshold = true;
509 reconcile_metric_sinks = true;
510 }
511 catalog::Op::RenameItem { id, .. } => {
512 let item = self.catalog().get_entry(id);
513 let is_webhook_source = item
514 .source()
515 .map(|s| matches!(s.data_source, DataSourceDesc::Webhook { .. }))
516 .unwrap_or(false);
517 if is_webhook_source {
518 webhook_sources_to_restart.insert(*id);
519 }
520 }
521 catalog::Op::RenameSchema {
522 database_spec,
523 schema_spec,
524 ..
525 } => {
526 let schema = self.catalog().get_schema(
527 database_spec,
528 schema_spec,
529 conn_id.unwrap_or(&SYSTEM_CONN_ID),
530 );
531 let webhook_sources = schema.item_ids().filter(|id| {
532 let item = self.catalog().get_entry(id);
533 item.source()
534 .map(|s| matches!(s.data_source, DataSourceDesc::Webhook { .. }))
535 .unwrap_or(false)
536 });
537 webhook_sources_to_restart.extend(webhook_sources);
538 }
539 catalog::Op::CreateCluster { id, .. } => {
540 clusters_to_create.push(*id);
541 }
542 catalog::Op::CreateClusterReplica {
543 cluster_id,
544 name,
545 config,
546 ..
547 } => {
548 cluster_replicas_to_create.push((
549 *cluster_id,
550 name.clone(),
551 config.location.num_processes(),
552 ));
553 }
554 _ => (),
555 }
556 }
557
558 let validate_res = self.validate_resource_limits(&ops, conn_id.unwrap_or(&SYSTEM_CONN_ID));
561 phase_seconds
562 .with_label_values(&["prep"])
563 .observe(phase_start.elapsed().as_secs_f64());
564 validate_res?;
565
566 let oracle_write_ts = self
575 .get_catalog_write_ts()
576 .wall_time()
577 .observe(phase_seconds.with_label_values(&["write_ts"]))
578 .await;
579
580 let Coordinator {
581 catalog,
582 active_conns,
583 controller,
584 cluster_replica_statuses,
585 ..
586 } = self;
587 let catalog = Arc::make_mut(catalog);
588 let conn = conn_id.map(|id| active_conns.get(id).expect("connection must exist"));
589
590 if let Some(conn) = conn {
593 let creates_temp_item = ops.iter().any(
594 |op| matches!(op, catalog::Op::CreateItem { item, .. } if item.is_temporary()),
595 );
596 if creates_temp_item && !catalog.state().has_temporary_namespace(conn.conn_id()) {
597 catalog.register_temporary_namespace(conn.conn_id(), conn.uuid());
598 }
599 }
600
601 let TransactionResult {
610 builtin_table_updates,
611 catalog_updates,
612 audit_events,
613 } = catalog
614 .transact(
615 Some(&mut controller.storage_collections),
616 oracle_write_ts,
617 conn,
618 ops,
619 )
620 .wall_time()
621 .observe(phase_seconds.with_label_values(&["transact"]))
622 .await?;
623
624 for (cluster_id, replica_id) in &cluster_replicas_to_drop {
625 cluster_replica_statuses.remove_cluster_replica_statuses(cluster_id, replica_id);
626 }
627 for cluster_id in &clusters_to_drop {
628 cluster_replica_statuses.remove_cluster_statuses(cluster_id);
629 }
630 for cluster_id in clusters_to_create {
631 cluster_replica_statuses.initialize_cluster_statuses(cluster_id);
632 }
633 let now = to_datetime((catalog.config().now)());
634 for (cluster_id, replica_name, num_processes) in cluster_replicas_to_create {
635 let replica_id = catalog
636 .resolve_replica_in_cluster(&cluster_id, &replica_name)
637 .expect("just created")
638 .replica_id();
639 cluster_replica_statuses.initialize_cluster_replica_statuses(
640 cluster_id,
641 replica_id,
642 num_processes,
643 now,
644 );
645 }
646
647 let stage_start = Instant::now();
650 let builtin_update_notify = self.builtin_table_update().execute(builtin_table_updates);
651 phase_seconds
652 .with_label_values(&["stage_builtin"])
653 .observe(stage_start.elapsed().as_secs_f64());
654
655 let finalize_start = Instant::now();
656
657 let _: () = async {
660 if !webhook_sources_to_restart.is_empty() {
661 self.restart_webhook_sources(webhook_sources_to_restart);
662 }
663
664 if update_metrics_config {
665 mz_metrics::update_dyncfg(&self.catalog().system_config().dyncfg_updates());
666 }
667 if update_controller_config {
668 self.update_controller_config();
669 }
670 if update_compute_config {
671 self.update_compute_config();
672 }
673 if update_storage_config {
674 self.update_storage_config();
675 }
676 if update_timestamp_oracle_config {
677 self.update_timestamp_oracle_config();
678 }
679 if update_metrics_retention {
680 self.update_metrics_retention();
681 }
682 if update_tracing_config {
683 self.update_tracing_config();
684 }
685 if update_secrets_caching_config {
686 self.update_secrets_caching_config();
687 }
688 if update_cluster_scheduling_config {
689 self.update_cluster_scheduling_config();
690 }
691 if update_http_config {
692 self.update_http_config();
693 }
694 if update_advance_timelines_interval {
695 let new_interval = self.catalog().system_config().default_timestamp_interval();
696 if new_interval != self.advance_timelines_interval.period() {
697 self.advance_timelines_interval = tokio::time::interval(new_interval);
698 }
699 }
700 if update_optimizer_e2e_latency_warning_threshold {
701 let threshold = self
702 .catalog()
703 .system_config()
704 .optimizer_e2e_latency_warning_threshold();
705 self.optimizer_metrics
706 .set_e2e_optimization_time_log_threshold(threshold);
707 }
708 if reconcile_metric_sinks {
709 self.reconcile_metric_sinks().await;
710 }
711 }
712 .instrument(info_span!("coord::catalog_transact_with::finalize"))
713 .await;
714
715 let conn = conn_id.and_then(|id| self.active_conns.get(id));
716 if let Some(segment_client) = &self.segment_client {
717 for VersionedEvent::V1(event) in audit_events {
718 let event_type = format!(
719 "{} {}",
720 event.object_type.as_title_case(),
721 event.event_type.as_title_case()
722 );
723 segment_client.environment_track(
724 &self.catalog().config().environment_id,
725 event_type,
726 json!({ "details": event.details.as_json() }),
727 EventDetails {
728 user_id: conn
729 .and_then(|c| c.user().external_metadata.as_ref())
730 .map(|m| m.user_id),
731 application_name: conn.map(|c| c.application_name()),
732 ..Default::default()
733 },
734 );
735 }
736 }
737
738 phase_seconds
739 .with_label_values(&["finalize"])
740 .observe(finalize_start.elapsed().as_secs_f64());
741
742 Ok((builtin_update_notify, catalog_updates))
743 }
744
745 pub(crate) fn drop_replica(&mut self, cluster_id: ClusterId, replica_id: ReplicaId) {
746 self.drop_introspection_subscribes(replica_id);
747 self.drop_metric_sinks(replica_id);
748
749 self.controller
750 .drop_replica(cluster_id, replica_id)
751 .expect("dropping replica must not fail");
752 }
753
754 pub(crate) fn drop_sources(&mut self, sources: Vec<(CatalogItemId, GlobalId)>) {
756 for (item_id, _gid) in &sources {
757 self.active_webhooks.remove(item_id);
758 }
759 let storage_metadata = self.catalog.state().storage_metadata();
760 let source_gids = sources.into_iter().map(|(_id, gid)| gid).collect();
761 self.controller
762 .storage
763 .drop_sources(storage_metadata, source_gids)
764 .unwrap_or_terminate("cannot fail to drop sources");
765 }
766
767 pub(crate) async fn drop_tables(&mut self, tables: Vec<(CatalogItemId, GlobalId)>) {
769 for (item_id, _gid) in &tables {
770 self.active_webhooks.remove(item_id);
771 }
772
773 let table_gids: Vec<_> = tables.into_iter().map(|(_id, gid)| gid).collect();
774
775 let forget_ids = self
777 .controller
778 .storage
779 .txns_table_ids(table_gids.clone())
780 .unwrap_or_terminate("cannot fail to look up txns-registered tables");
781 if !forget_ids.is_empty() {
782 self.forget_tables_via_committer(forget_ids).await;
783 }
784
785 let storage_metadata = self.catalog.state().storage_metadata();
786 self.controller
787 .storage
788 .drop_tables(storage_metadata, table_gids)
789 .unwrap_or_terminate("cannot fail to drop tables");
790 }
791
792 fn restart_webhook_sources(&mut self, sources: impl IntoIterator<Item = CatalogItemId>) {
793 for id in sources {
794 self.active_webhooks.remove(&id);
795 }
796 }
797
798 #[must_use]
804 pub async fn drop_compute_sink(
805 &mut self,
806 sink_id: GlobalId,
807 ) -> Option<(ActiveComputeSink, BuiltinTableAppendNotify)> {
808 self.drop_compute_sinks([sink_id]).await.remove(&sink_id)
809 }
810
811 #[must_use]
822 pub async fn drop_compute_sinks(
823 &mut self,
824 sink_ids: impl IntoIterator<Item = GlobalId>,
825 ) -> BTreeMap<GlobalId, (ActiveComputeSink, BuiltinTableAppendNotify)> {
826 let mut by_id = BTreeMap::new();
827 let mut by_cluster: BTreeMap<_, Vec<_>> = BTreeMap::new();
828 for sink_id in sink_ids {
829 let (sink, write_notify) = match self.remove_active_compute_sink(sink_id).await {
830 None => {
831 tracing::debug!(%sink_id, "drop_compute_sinks: sink already removed");
836 continue;
837 }
838 Some(entry) => entry,
839 };
840
841 by_cluster
842 .entry(sink.cluster_id())
843 .or_default()
844 .push(sink_id);
845 by_id.insert(sink_id, (sink, write_notify));
846 }
847 for (cluster_id, ids) in by_cluster {
848 let compute = &mut self.controller.compute;
849 if compute.instance_exists(cluster_id) {
851 compute
852 .drop_collections(cluster_id, ids)
853 .unwrap_or_terminate("cannot fail to drop collections");
854 }
855 }
856 by_id
857 }
858
859 pub async fn retire_compute_sinks(
865 &mut self,
866 mut reasons: BTreeMap<GlobalId, ActiveComputeSinkRetireReason>,
867 ) -> BuiltinTableAppendCompletion {
868 let sink_ids = reasons.keys().cloned();
869 let to_retire: Vec<_> = self
870 .drop_compute_sinks(sink_ids)
871 .await
872 .into_iter()
873 .map(|(id, (sink, write_notify))| {
874 let reason = reasons
875 .remove(&id)
876 .expect("all returned IDs are in `reasons`");
877 (sink, write_notify, reason)
878 })
879 .collect();
880
881 let (done_tx, done_rx) = tokio::sync::oneshot::channel();
887 task::spawn(|| "retire_compute_sinks", async move {
888 for (sink, write_notify, reason) in to_retire {
889 write_notify.await;
890 sink.retire(reason);
891 }
892 let _ = done_tx.send(());
893 });
894 BuiltinTableAppendCompletion::new(Box::pin(async move {
895 let _ = done_rx.await;
896 }))
897 }
898
899 #[mz_ore::instrument(level = "debug")]
901 pub(crate) async fn cancel_compute_sinks_for_conn(
902 &mut self,
903 conn_id: &ConnectionId,
904 ) -> BuiltinTableAppendCompletion {
905 self.retire_compute_sinks_for_conn(conn_id, ActiveComputeSinkRetireReason::Canceled)
906 .await
907 }
908
909 #[mz_ore::instrument(level = "debug")]
912 pub(crate) async fn retire_compute_sinks_for_conn(
913 &mut self,
914 conn_id: &ConnectionId,
915 reason: ActiveComputeSinkRetireReason,
916 ) -> BuiltinTableAppendCompletion {
917 let drop_sinks = self
918 .active_conns
919 .get_mut(conn_id)
920 .expect("must exist for active session")
921 .drop_sinks
922 .iter()
923 .map(|sink_id| (*sink_id, reason.clone()))
924 .collect();
925 self.retire_compute_sinks(drop_sinks).await
926 }
927
928 pub(crate) fn drop_storage_sinks(&mut self, sink_gids: Vec<GlobalId>) {
929 let storage_metadata = self.catalog.state().storage_metadata();
930 self.controller
931 .storage
932 .drop_sinks(storage_metadata, sink_gids)
933 .unwrap_or_terminate("cannot fail to drop sinks");
934 }
935
936 pub(crate) fn drop_compute_collections(&mut self, collections: Vec<(ClusterId, GlobalId)>) {
937 let mut by_cluster: BTreeMap<_, Vec<_>> = BTreeMap::new();
938 for (cluster_id, gid) in collections {
939 by_cluster.entry(cluster_id).or_default().push(gid);
940 }
941 for (cluster_id, gids) in by_cluster {
942 let compute = &mut self.controller.compute;
943 if compute.instance_exists(cluster_id) {
945 compute
946 .drop_collections(cluster_id, gids)
947 .unwrap_or_terminate("cannot fail to drop collections");
948 }
949 }
950 }
951
952 pub(crate) fn drop_vpc_endpoints_in_background(&self, vpc_endpoints: Vec<CatalogItemId>) {
953 let Some(cloud_resource_controller) = self.cloud_resource_controller.as_ref() else {
957 warn!("dropping VPC endpoints without cloud_resource_controller; skipping cleanup");
958 return;
959 };
960 let cloud_resource_controller = Arc::clone(cloud_resource_controller);
961 task::spawn(
969 || "drop_vpc_endpoints",
970 async move {
971 for vpc_endpoint in vpc_endpoints {
972 let _ = Retry::default()
973 .max_duration(Duration::from_secs(60))
974 .retry_async(|_state| async {
975 fail_point!("drop_vpc_endpoint", |r| {
976 Err(anyhow::anyhow!("Fail point error {:?}", r))
977 });
978 match cloud_resource_controller
979 .delete_vpc_endpoint(vpc_endpoint)
980 .await
981 {
982 Ok(_) => Ok(()),
983 Err(e) => {
984 warn!("Dropping VPC Endpoints has encountered an error: {}", e);
985 Err(e)
986 }
987 }
988 })
989 .await;
990 }
991 }
992 .instrument(info_span!(
993 "coord::catalog_transact_inner::drop_vpc_endpoints"
994 )),
995 );
996 }
997
998 pub(crate) async fn drop_temp_items(&mut self, conn_id: &ConnectionId) {
1001 let temp_items = self.catalog().state().get_temp_items(conn_id).collect();
1002 let all_items = self.catalog().object_dependents(&temp_items, conn_id);
1003
1004 if all_items.is_empty() {
1005 return;
1006 }
1007 let op = Op::DropObjects(
1008 all_items
1009 .into_iter()
1010 .map(DropObjectInfo::manual_drop_from_object_id)
1011 .collect(),
1012 );
1013
1014 self.catalog_transact_with_context(Some(conn_id), None, vec![op])
1015 .await
1016 .expect("unable to drop temporary items for conn_id");
1017 }
1018
1019 fn update_cluster_scheduling_config(&self) {
1020 let config = flags::orchestrator_scheduling_config(self.catalog.system_config());
1021 self.controller
1022 .update_orchestrator_scheduling_config(config);
1023 }
1024
1025 fn update_secrets_caching_config(&self) {
1026 let config = flags::caching_config(self.catalog.system_config());
1027 self.caching_secrets_reader.set_policy(config);
1028 }
1029
1030 fn update_tracing_config(&self) {
1031 let tracing = flags::tracing_config(self.catalog().system_config());
1032 tracing.apply(&self.tracing_handle);
1033 }
1034
1035 fn update_compute_config(&mut self) {
1036 let config_params = flags::compute_config(self.catalog().system_config());
1037 self.controller.compute.update_configuration(config_params);
1038 }
1039
1040 fn update_storage_config(&mut self) {
1041 let config_params = flags::storage_config(self.catalog().system_config());
1042 self.controller.storage.update_parameters(config_params);
1043 }
1044
1045 fn update_timestamp_oracle_config(&self) {
1046 let config_params = flags::timestamp_oracle_config(self.catalog().system_config());
1047 if let Some(config) = self.timestamp_oracle_config.as_ref() {
1048 config.apply_parameters(config_params)
1049 }
1050 }
1051
1052 fn update_metrics_retention(&self) {
1053 let duration = self.catalog().system_config().metrics_retention();
1054 let policy = ReadPolicy::lag_writes_by(
1055 Timestamp::new(u64::try_from(duration.as_millis()).unwrap_or_else(|_e| {
1056 tracing::error!("Absurd metrics retention duration: {duration:?}.");
1057 u64::MAX
1058 })),
1059 SINCE_GRANULARITY,
1060 );
1061 let storage_policies = self
1062 .catalog()
1063 .entries()
1064 .filter(|entry| {
1065 entry.item().is_retained_metrics_object()
1066 && entry.item().is_compute_object_on_cluster().is_none()
1067 })
1068 .map(|entry| (entry.id(), policy.clone()))
1069 .collect::<Vec<_>>();
1070 let compute_policies = self
1071 .catalog()
1072 .entries()
1073 .filter_map(|entry| {
1074 if let (true, Some(cluster_id)) = (
1075 entry.item().is_retained_metrics_object(),
1076 entry.item().is_compute_object_on_cluster(),
1077 ) {
1078 Some((cluster_id, entry.id(), policy.clone()))
1079 } else {
1080 None
1081 }
1082 })
1083 .collect::<Vec<_>>();
1084 self.update_storage_read_policies(storage_policies);
1085 self.update_compute_read_policies(compute_policies);
1086 }
1087
1088 fn update_controller_config(&mut self) {
1089 let sys_config = self.catalog().system_config();
1090 self.controller
1091 .update_configuration(sys_config.dyncfg_updates());
1092 }
1093
1094 fn update_http_config(&mut self) {
1095 let webhook_request_limit = self
1096 .catalog()
1097 .system_config()
1098 .webhook_concurrent_request_limit();
1099 self.webhook_concurrency_limit
1100 .set_limit(webhook_request_limit);
1101 }
1102
1103 pub(crate) async fn create_storage_export(
1104 &mut self,
1105 id: GlobalId,
1106 sink: &Sink,
1107 ) -> Result<(), AdapterError> {
1108 self.controller.storage.check_exists(sink.from)?;
1110
1111 let id_bundle = crate::CollectionIdBundle {
1118 storage_ids: btreeset! {sink.from},
1119 compute_ids: btreemap! {},
1120 };
1121
1122 let read_holds = self.acquire_read_holds(&id_bundle);
1130 let as_of = read_holds.least_valid_read();
1131
1132 let storage_sink_from_entry = self.catalog().get_entry_by_global_id(&sink.from);
1133 let storage_sink_desc = mz_storage_types::sinks::StorageSinkDesc {
1134 from: sink.from,
1135 from_desc: storage_sink_from_entry
1136 .relation_desc()
1137 .expect("sinks can only be built on items with descs")
1138 .into_owned(),
1139 connection: sink
1140 .connection
1141 .clone()
1142 .into_inline_connection(self.catalog().state()),
1143 envelope: sink.envelope,
1144 as_of,
1145 with_snapshot: sink.with_snapshot,
1146 version: sink.version,
1147 from_storage_metadata: (),
1148 to_storage_metadata: (),
1149 commit_interval: sink.commit_interval,
1150 };
1151
1152 let collection_desc = CollectionDescription {
1153 desc: KAFKA_PROGRESS_DESC.clone(),
1155 data_source: DataSource::Sink {
1156 desc: ExportDescription {
1157 sink: storage_sink_desc,
1158 instance_id: sink.cluster_id,
1159 },
1160 },
1161 since: None,
1162 timeline: None,
1163 primary: None,
1164 };
1165 let collections = vec![(id, collection_desc)];
1166
1167 let storage_metadata = self.catalog.state().storage_metadata();
1169 let res = self
1170 .controller
1171 .storage
1172 .create_collections(storage_metadata, None, collections)
1173 .await;
1174
1175 drop(read_holds);
1178
1179 Ok(res?)
1180 }
1181
1182 fn validate_resource_limits(
1185 &self,
1186 ops: &Vec<catalog::Op>,
1187 conn_id: &ConnectionId,
1188 ) -> Result<(), AdapterError> {
1189 let mut new_kafka_connections = 0;
1190 let mut new_postgres_connections = 0;
1191 let mut new_mysql_connections = 0;
1192 let mut new_sql_server_connections = 0;
1193 let mut new_aws_privatelink_connections = 0;
1194 let mut new_tables = 0;
1195 let mut new_sources = 0;
1196 let mut new_sinks = 0;
1197 let mut new_materialized_views = 0;
1198 let mut new_clusters = 0;
1199 let mut new_replicas_per_cluster = BTreeMap::new();
1200 let mut new_credit_consumption_rate = Numeric::zero();
1201 let mut new_databases = 0;
1202 let mut new_schemas_per_database = BTreeMap::new();
1203 let mut new_objects_per_schema = BTreeMap::new();
1204 let mut new_secrets = 0;
1205 let mut new_roles = 0;
1206 let mut new_network_policies = 0;
1207 for op in ops {
1208 match op {
1209 Op::CreateDatabase { .. } => {
1210 new_databases += 1;
1211 }
1212 Op::CreateSchema { database_id, .. } => {
1213 if let ResolvedDatabaseSpecifier::Id(database_id) = database_id {
1214 *new_schemas_per_database.entry(database_id).or_insert(0) += 1;
1215 }
1216 }
1217 Op::CreateRole { .. } => {
1218 new_roles += 1;
1219 }
1220 Op::CreateNetworkPolicy { .. } => {
1221 new_network_policies += 1;
1222 }
1223 Op::CreateCluster { .. } => {
1224 new_clusters += 1;
1228 }
1229 Op::CreateClusterReplica {
1230 cluster_id, config, ..
1231 } => {
1232 if cluster_id.is_user() {
1233 *new_replicas_per_cluster.entry(*cluster_id).or_insert(0) += 1;
1234 if let ReplicaLocation::Managed(location) = &config.location {
1235 new_credit_consumption_rate += self.replica_credits_per_hour(location);
1236 }
1237 }
1238 }
1239 Op::CreateItem { name, item, .. } => {
1240 *new_objects_per_schema
1241 .entry((
1242 name.qualifiers.database_spec.clone(),
1243 name.qualifiers.schema_spec.clone(),
1244 ))
1245 .or_insert(0) += 1;
1246 match item {
1247 CatalogItem::Connection(connection) => match connection.details {
1248 ConnectionDetails::Kafka(_) => new_kafka_connections += 1,
1249 ConnectionDetails::Postgres(_) => new_postgres_connections += 1,
1250 ConnectionDetails::MySql(_) => new_mysql_connections += 1,
1251 ConnectionDetails::SqlServer(_) => new_sql_server_connections += 1,
1252 ConnectionDetails::AwsPrivatelink(_) => {
1253 new_aws_privatelink_connections += 1
1254 }
1255 ConnectionDetails::Csr(_)
1256 | ConnectionDetails::GlueSchemaRegistry(_)
1257 | ConnectionDetails::Ssh { .. }
1258 | ConnectionDetails::Aws(_)
1259 | ConnectionDetails::Gcp(_)
1260 | ConnectionDetails::IcebergCatalog(_) => {}
1261 },
1262 CatalogItem::Table(_) => {
1263 new_tables += 1;
1264 }
1265 CatalogItem::Source(source) => {
1266 new_sources += source.user_controllable_persist_shard_count()
1267 }
1268 CatalogItem::Sink(_) => new_sinks += 1,
1269 CatalogItem::MaterializedView(_) => {
1270 new_materialized_views += 1;
1271 }
1272 CatalogItem::Secret(_) => {
1273 new_secrets += 1;
1274 }
1275 CatalogItem::Log(_)
1276 | CatalogItem::View(_)
1277 | CatalogItem::Index(_)
1278 | CatalogItem::Type(_)
1279 | CatalogItem::Func(_)
1280 | CatalogItem::MetricSink(_) => {}
1281 }
1282 }
1283 Op::DropObjects(drop_object_infos) => {
1284 for drop_object_info in drop_object_infos {
1285 match drop_object_info {
1286 DropObjectInfo::Cluster(_) => {
1287 new_clusters -= 1;
1288 }
1289 DropObjectInfo::ClusterReplica((cluster_id, replica_id, _reason)) => {
1290 if cluster_id.is_user() {
1291 *new_replicas_per_cluster.entry(*cluster_id).or_insert(0) -= 1;
1292 let cluster = self
1293 .catalog()
1294 .get_cluster_replica(*cluster_id, *replica_id);
1295 if let ReplicaLocation::Managed(location) =
1296 &cluster.config.location
1297 {
1298 new_credit_consumption_rate -=
1299 self.replica_credits_per_hour(location);
1300 }
1301 }
1302 }
1303 DropObjectInfo::Database(_) => {
1304 new_databases -= 1;
1305 }
1306 DropObjectInfo::Schema((database_spec, _)) => {
1307 if let ResolvedDatabaseSpecifier::Id(database_id) = database_spec {
1308 *new_schemas_per_database.entry(database_id).or_insert(0) -= 1;
1309 }
1310 }
1311 DropObjectInfo::Role(_) => {
1312 new_roles -= 1;
1313 }
1314 DropObjectInfo::NetworkPolicy(_) => {
1315 new_network_policies -= 1;
1316 }
1317 DropObjectInfo::Item(id) => {
1318 let entry = self.catalog().get_entry(id);
1319 *new_objects_per_schema
1320 .entry((
1321 entry.name().qualifiers.database_spec.clone(),
1322 entry.name().qualifiers.schema_spec.clone(),
1323 ))
1324 .or_insert(0) -= 1;
1325 match entry.item() {
1326 CatalogItem::Connection(connection) => match connection.details
1327 {
1328 ConnectionDetails::AwsPrivatelink(_) => {
1329 new_aws_privatelink_connections -= 1;
1330 }
1331 _ => (),
1332 },
1333 CatalogItem::Table(_) => {
1334 new_tables -= 1;
1335 }
1336 CatalogItem::Source(source) => {
1337 new_sources -=
1338 source.user_controllable_persist_shard_count()
1339 }
1340 CatalogItem::Sink(_) => new_sinks -= 1,
1341 CatalogItem::MaterializedView(_) => {
1342 new_materialized_views -= 1;
1343 }
1344 CatalogItem::Secret(_) => {
1345 new_secrets -= 1;
1346 }
1347 CatalogItem::Log(_)
1348 | CatalogItem::View(_)
1349 | CatalogItem::Index(_)
1350 | CatalogItem::Type(_)
1351 | CatalogItem::Func(_)
1352 | CatalogItem::MetricSink(_) => {}
1353 }
1354 }
1355 }
1356 }
1357 }
1358 Op::UpdateItem {
1359 name: _,
1360 id,
1361 to_item,
1362 } => match to_item {
1363 CatalogItem::Source(source) => {
1364 let current_source = self
1365 .catalog()
1366 .get_entry(id)
1367 .source()
1368 .expect("source update is for source item");
1369
1370 new_sources += source.user_controllable_persist_shard_count()
1371 - current_source.user_controllable_persist_shard_count();
1372 }
1373 CatalogItem::Connection(_)
1374 | CatalogItem::Table(_)
1375 | CatalogItem::Sink(_)
1376 | CatalogItem::MaterializedView(_)
1377 | CatalogItem::Secret(_)
1378 | CatalogItem::Log(_)
1379 | CatalogItem::View(_)
1380 | CatalogItem::Index(_)
1381 | CatalogItem::Type(_)
1382 | CatalogItem::Func(_)
1383 | CatalogItem::MetricSink(_) => {}
1384 },
1385 Op::AlterRole { .. }
1386 | Op::AlterRetainHistory { .. }
1387 | Op::AlterSourceTimestampInterval { .. }
1388 | Op::AlterNetworkPolicy { .. }
1389 | Op::AlterAddColumn { .. }
1390 | Op::AlterMaterializedViewApplyReplacement { .. }
1391 | Op::UpdatePrivilege { .. }
1392 | Op::UpdateDefaultPrivilege { .. }
1393 | Op::GrantRole { .. }
1394 | Op::RenameCluster { .. }
1395 | Op::RenameClusterReplica { .. }
1396 | Op::RenameItem { .. }
1397 | Op::RenameSchema { .. }
1398 | Op::UpdateOwner { .. }
1399 | Op::RevokeRole { .. }
1400 | Op::UpdateClusterConfig { .. }
1401 | Op::UpdateSourceReferences { .. }
1402 | Op::UpdateSystemConfiguration { .. }
1403 | Op::ResetSystemConfiguration { .. }
1404 | Op::ResetAllSystemConfiguration { .. }
1405 | Op::UpdateScopedSystemParameters { .. }
1406 | Op::Comment { .. }
1407 | Op::CheckClusterState { .. }
1408 | Op::InjectAuditEvents { .. } => {}
1409 }
1410 }
1411
1412 let mut current_aws_privatelink_connections = 0;
1413 let mut current_postgres_connections = 0;
1414 let mut current_mysql_connections = 0;
1415 let mut current_sql_server_connections = 0;
1416 let mut current_kafka_connections = 0;
1417 for c in self.catalog().user_connections() {
1418 let connection = c
1419 .connection()
1420 .expect("`user_connections()` only returns connection objects");
1421
1422 match connection.details {
1423 ConnectionDetails::AwsPrivatelink(_) => current_aws_privatelink_connections += 1,
1424 ConnectionDetails::Postgres(_) => current_postgres_connections += 1,
1425 ConnectionDetails::MySql(_) => current_mysql_connections += 1,
1426 ConnectionDetails::SqlServer(_) => current_sql_server_connections += 1,
1427 ConnectionDetails::Kafka(_) => current_kafka_connections += 1,
1428 ConnectionDetails::Csr(_)
1429 | ConnectionDetails::GlueSchemaRegistry(_)
1430 | ConnectionDetails::Ssh { .. }
1431 | ConnectionDetails::Aws(_)
1432 | ConnectionDetails::Gcp(_)
1433 | ConnectionDetails::IcebergCatalog(_) => {}
1434 }
1435 }
1436 self.validate_resource_limit(
1437 current_kafka_connections,
1438 new_kafka_connections,
1439 SystemVars::max_kafka_connections,
1440 "Kafka Connection",
1441 MAX_KAFKA_CONNECTIONS.name(),
1442 )?;
1443 self.validate_resource_limit(
1444 current_postgres_connections,
1445 new_postgres_connections,
1446 SystemVars::max_postgres_connections,
1447 "PostgreSQL Connection",
1448 MAX_POSTGRES_CONNECTIONS.name(),
1449 )?;
1450 self.validate_resource_limit(
1451 current_mysql_connections,
1452 new_mysql_connections,
1453 SystemVars::max_mysql_connections,
1454 "MySQL Connection",
1455 MAX_MYSQL_CONNECTIONS.name(),
1456 )?;
1457 self.validate_resource_limit(
1458 current_sql_server_connections,
1459 new_sql_server_connections,
1460 SystemVars::max_sql_server_connections,
1461 "SQL Server Connection",
1462 MAX_SQL_SERVER_CONNECTIONS.name(),
1463 )?;
1464 self.validate_resource_limit(
1465 current_aws_privatelink_connections,
1466 new_aws_privatelink_connections,
1467 SystemVars::max_aws_privatelink_connections,
1468 "AWS PrivateLink Connection",
1469 MAX_AWS_PRIVATELINK_CONNECTIONS.name(),
1470 )?;
1471 self.validate_resource_limit(
1472 self.catalog().user_tables().count(),
1473 new_tables,
1474 SystemVars::max_tables,
1475 "table",
1476 MAX_TABLES.name(),
1477 )?;
1478
1479 let current_sources: usize = self
1480 .catalog()
1481 .user_sources()
1482 .filter_map(|source| source.source())
1483 .map(|source| source.user_controllable_persist_shard_count())
1484 .sum::<i64>()
1485 .try_into()
1486 .expect("non-negative sum of sources");
1487
1488 self.validate_resource_limit(
1489 current_sources,
1490 new_sources,
1491 SystemVars::max_sources,
1492 "source",
1493 MAX_SOURCES.name(),
1494 )?;
1495 self.validate_resource_limit(
1496 self.catalog().user_sinks().count(),
1497 new_sinks,
1498 SystemVars::max_sinks,
1499 "sink",
1500 MAX_SINKS.name(),
1501 )?;
1502 self.validate_resource_limit(
1503 self.catalog().user_materialized_views().count(),
1504 new_materialized_views,
1505 SystemVars::max_materialized_views,
1506 "materialized view",
1507 MAX_MATERIALIZED_VIEWS.name(),
1508 )?;
1509 self.validate_resource_limit(
1510 self.catalog().user_clusters().count(),
1516 new_clusters,
1517 SystemVars::max_clusters,
1518 "cluster",
1519 MAX_CLUSTERS.name(),
1520 )?;
1521 for (cluster_id, new_replicas) in new_replicas_per_cluster {
1522 let current_amount = self
1524 .catalog()
1525 .try_get_cluster(cluster_id)
1526 .map(|instance| instance.user_replicas().count())
1527 .unwrap_or(0);
1528 self.validate_resource_limit(
1529 current_amount,
1530 new_replicas,
1531 SystemVars::max_replicas_per_cluster,
1532 "cluster replica",
1533 MAX_REPLICAS_PER_CLUSTER.name(),
1534 )?;
1535 }
1536 self.validate_resource_limit_numeric(
1537 self.current_credit_consumption_rate(None),
1538 new_credit_consumption_rate,
1539 |system_vars| {
1540 self.license_key
1541 .max_credit_consumption_rate()
1542 .map_or_else(|| system_vars.max_credit_consumption_rate(), Numeric::from)
1543 },
1544 "cluster replica",
1545 MAX_CREDIT_CONSUMPTION_RATE.name(),
1546 )?;
1547 self.validate_resource_limit(
1548 self.catalog().databases().count(),
1549 new_databases,
1550 SystemVars::max_databases,
1551 "database",
1552 MAX_DATABASES.name(),
1553 )?;
1554 for (database_id, new_schemas) in new_schemas_per_database {
1555 self.validate_resource_limit(
1556 self.catalog().get_database(database_id).schemas_by_id.len(),
1557 new_schemas,
1558 SystemVars::max_schemas_per_database,
1559 "schema",
1560 MAX_SCHEMAS_PER_DATABASE.name(),
1561 )?;
1562 }
1563 for ((database_spec, schema_spec), new_objects) in new_objects_per_schema {
1564 let current_items = self
1567 .catalog()
1568 .try_get_schema(&database_spec, &schema_spec, conn_id)
1569 .map(|schema| schema.items.len())
1570 .unwrap_or(0);
1571 self.validate_resource_limit(
1572 current_items,
1573 new_objects,
1574 SystemVars::max_objects_per_schema,
1575 "object",
1576 MAX_OBJECTS_PER_SCHEMA.name(),
1577 )?;
1578 }
1579 self.validate_resource_limit(
1580 self.catalog().user_secrets().count(),
1581 new_secrets,
1582 SystemVars::max_secrets,
1583 "secret",
1584 MAX_SECRETS.name(),
1585 )?;
1586 self.validate_resource_limit(
1587 self.catalog().user_roles().count(),
1588 new_roles,
1589 SystemVars::max_roles,
1590 "role",
1591 MAX_ROLES.name(),
1592 )?;
1593 self.validate_resource_limit(
1594 self.catalog().user_network_policies().count(),
1595 new_network_policies,
1596 SystemVars::max_network_policies,
1597 "network_policy",
1598 MAX_NETWORK_POLICIES.name(),
1599 )?;
1600 Ok(())
1601 }
1602
1603 pub(crate) fn validate_resource_limit<F>(
1605 &self,
1606 current_amount: usize,
1607 new_instances: i64,
1608 resource_limit: F,
1609 resource_type: &str,
1610 limit_name: &str,
1611 ) -> Result<(), AdapterError>
1612 where
1613 F: Fn(&SystemVars) -> u32,
1614 {
1615 if new_instances <= 0 {
1616 return Ok(());
1617 }
1618
1619 let limit: i64 = resource_limit(self.catalog().system_config()).into();
1620 let current_amount: Option<i64> = current_amount.try_into().ok();
1621 let desired =
1622 current_amount.and_then(|current_amount| current_amount.checked_add(new_instances));
1623
1624 let exceeds_limit = if let Some(desired) = desired {
1625 desired > limit
1626 } else {
1627 true
1628 };
1629
1630 let desired = desired
1631 .map(|desired| desired.to_string())
1632 .unwrap_or_else(|| format!("more than {}", i64::MAX));
1633 let current = current_amount
1634 .map(|current| current.to_string())
1635 .unwrap_or_else(|| format!("more than {}", i64::MAX));
1636 if exceeds_limit {
1637 Err(AdapterError::ResourceExhaustion {
1638 resource_type: resource_type.to_string(),
1639 limit_name: limit_name.to_string(),
1640 desired,
1641 limit: limit.to_string(),
1642 current,
1643 })
1644 } else {
1645 Ok(())
1646 }
1647 }
1648
1649 pub(crate) fn validate_resource_limit_numeric<F>(
1653 &self,
1654 current_amount: Numeric,
1655 new_amount: Numeric,
1656 resource_limit: F,
1657 resource_type: &str,
1658 limit_name: &str,
1659 ) -> Result<(), AdapterError>
1660 where
1661 F: Fn(&SystemVars) -> Numeric,
1662 {
1663 if new_amount <= Numeric::zero() {
1664 return Ok(());
1665 }
1666
1667 let limit = resource_limit(self.catalog().system_config());
1668 let desired = current_amount + new_amount;
1672 if desired > limit {
1673 Err(AdapterError::ResourceExhaustion {
1674 resource_type: resource_type.to_string(),
1675 limit_name: limit_name.to_string(),
1676 desired: desired.to_string(),
1677 limit: limit.to_string(),
1678 current: current_amount.to_string(),
1679 })
1680 } else {
1681 Ok(())
1682 }
1683 }
1684}