1use std::collections::{BTreeMap, BTreeSet, VecDeque};
14use std::fmt::Debug;
15use std::iter;
16use std::str::FromStr;
17use std::sync::Arc;
18
19use differential_dataflow::consolidation::consolidate_updates;
20use futures::future;
21use itertools::{Either, Itertools};
22use mz_adapter_types::connection::ConnectionId;
23use mz_catalog::SYSTEM_CONN_ID;
24use mz_catalog::builtin::{
25 BUILTIN_LOG_LOOKUP, BUILTIN_LOOKUP, Builtin, BuiltinLog, BuiltinTable, BuiltinView,
26};
27use mz_catalog::durable::objects::{
28 ClusterKey, DatabaseKey, DurableType, ItemKey, NetworkPolicyKey, RoleAuthKey, RoleKey,
29 SchemaKey,
30};
31use mz_catalog::durable::{CatalogError, SystemObjectMapping};
32use mz_catalog::memory::error::{Error, ErrorKind};
33use mz_catalog::memory::objects::{
34 CatalogEntry, CatalogItem, Cluster, ClusterReplica, Database, Func, Index, Log, NetworkPolicy,
35 Role, RoleAuth, Schema, Source, StateDiff, StateUpdate, StateUpdateKind, Table,
36 TableDataSource, Type, UpdateFrom,
37};
38use mz_compute_types::config::ComputeReplicaConfig;
39use mz_compute_types::dataflows::DataflowDescription;
40use mz_controller::clusters::{ReplicaConfig, ReplicaLogging};
41use mz_controller_types::ClusterId;
42use mz_expr::MirScalarExpr;
43use mz_ore::collections::CollectionExt;
44use mz_ore::tracing::OpenTelemetryContext;
45use mz_ore::{
46 instrument, soft_assert_eq_or_log, soft_assert_no_log, soft_assert_or_log, soft_panic_or_log,
47};
48use mz_pgrepr::oid::INVALID_OID;
49use mz_repr::adt::mz_acl_item::{MzAclItem, PrivilegeMap};
50use mz_repr::role_id::RoleId;
51use mz_repr::{CatalogItemId, Diff, GlobalId, RelationVersion, Timestamp, VersionedRelationDesc};
52use mz_sql::catalog::CatalogError as SqlCatalogError;
53use mz_sql::catalog::{CatalogItem as SqlCatalogItem, CatalogItemType, CatalogSchema, CatalogType};
54use mz_sql::names::{
55 FullItemName, ItemQualifiers, QualifiedItemName, RawDatabaseSpecifier,
56 ResolvedDatabaseSpecifier, ResolvedIds, SchemaSpecifier,
57};
58use mz_sql::session::user::MZ_SYSTEM_ROLE_ID;
59use mz_sql::session::vars::{VarError, VarInput};
60use mz_sql::{plan, rbac};
61use mz_sql_parser::ast::Expr;
62use mz_storage_types::sources::Timeline;
63use mz_transform::dataflow::DataflowMetainfo;
64use mz_transform::notice::OptimizerNotice;
65use tracing::{info_span, warn};
66
67use crate::AdapterError;
68use crate::catalog::state::LocalExpressionCache;
69use crate::catalog::{BuiltinTableUpdate, CatalogState};
70use crate::coord::catalog_implications::parsed_state_updates::{self, ParsedStateUpdate};
71use crate::util::{index_sql, sort_topological};
72
73#[derive(Debug, Clone, Default)]
85struct InProgressRetractions {
86 roles: BTreeMap<RoleKey, Role>,
87 role_auths: BTreeMap<RoleAuthKey, RoleAuth>,
88 databases: BTreeMap<DatabaseKey, Database>,
89 schemas: BTreeMap<SchemaKey, Schema>,
90 clusters: BTreeMap<ClusterKey, Cluster>,
91 network_policies: BTreeMap<NetworkPolicyKey, NetworkPolicy>,
92 items: BTreeMap<ItemKey, CatalogEntry>,
93 introspection_source_indexes: BTreeMap<CatalogItemId, CatalogEntry>,
94 system_object_mappings: BTreeMap<CatalogItemId, CatalogEntry>,
95}
96
97impl CatalogState {
98 #[must_use]
102 #[instrument]
103 pub(crate) async fn apply_updates(
104 &mut self,
105 updates: Vec<StateUpdate>,
106 local_expression_cache: &mut LocalExpressionCache,
107 ) -> (
108 Vec<BuiltinTableUpdate<&'static BuiltinTable>>,
109 Vec<ParsedStateUpdate>,
110 ) {
111 let mut builtin_table_updates = Vec::with_capacity(updates.len());
112 let mut catalog_updates = Vec::with_capacity(updates.len());
113
114 let updates = Self::consolidate_updates(updates);
119
120 let mut groups: Vec<Vec<_>> = Vec::new();
122 for (_, updates) in &updates.into_iter().chunk_by(|update| update.ts) {
123 let updates = sort_updates(updates.collect());
127 groups.push(updates);
128 }
129
130 for updates in groups {
131 let mut apply_state = ApplyState::Updates(Vec::new());
132 let mut retractions = InProgressRetractions::default();
133
134 for update in updates {
135 let (next_apply_state, (builtin_table_update, catalog_update)) = apply_state
136 .step(
137 ApplyState::new(update),
138 self,
139 &mut retractions,
140 local_expression_cache,
141 )
142 .await;
143 apply_state = next_apply_state;
144 builtin_table_updates.extend(builtin_table_update);
145 catalog_updates.extend(catalog_update);
146 }
147
148 let (builtin_table_update, catalog_update) = apply_state
150 .apply(self, &mut retractions, local_expression_cache)
151 .await;
152 builtin_table_updates.extend(builtin_table_update);
153 catalog_updates.extend(catalog_update);
154
155 let dropped_entries: Vec<CatalogEntry> = retractions.items.into_values().collect();
158 if !dropped_entries.is_empty() {
159 let dropped_notices = self.drop_optimizer_notices(dropped_entries);
160 if self.system_config().enable_mz_notices() {
161 self.pack_optimizer_notice_updates(
162 &mut builtin_table_updates,
163 dropped_notices.iter(),
164 Diff::MINUS_ONE,
165 );
166 }
167 }
168 }
169
170 (builtin_table_updates, catalog_updates)
171 }
172
173 fn consolidate_updates(updates: Vec<StateUpdate>) -> Vec<StateUpdate> {
181 let mut updates: Vec<(StateUpdateKind, Timestamp, mz_repr::Diff)> = updates
182 .into_iter()
183 .map(|update| (update.kind, update.ts, update.diff.into()))
184 .collect_vec();
185
186 consolidate_updates(&mut updates);
187
188 updates
189 .into_iter()
190 .map(|(kind, ts, diff)| StateUpdate {
191 kind,
192 ts,
193 diff: diff
194 .try_into()
195 .expect("catalog state cannot have diff other than -1 or 1"),
196 })
197 .collect_vec()
198 }
199
200 #[instrument(level = "debug")]
201 fn apply_updates_inner(
202 &mut self,
203 updates: Vec<StateUpdate>,
204 retractions: &mut InProgressRetractions,
205 local_expression_cache: &mut LocalExpressionCache,
206 ) -> Result<
207 (
208 Vec<BuiltinTableUpdate<&'static BuiltinTable>>,
209 Vec<ParsedStateUpdate>,
210 ),
211 CatalogError,
212 > {
213 soft_assert_no_log!(
214 updates.iter().map(|update| update.ts).all_equal(),
215 "all timestamps should be equal: {updates:?}"
216 );
217
218 let mut update_system_config = false;
219
220 let mut builtin_table_updates = Vec::with_capacity(updates.len());
221 let mut catalog_updates = Vec::new();
222
223 for state_update in updates {
224 if self.is_nonlocal_ephemeral_item_update(&state_update) {
228 tracing::debug!(
229 ?state_update,
230 "skipping ephemeral item update for non-local session"
231 );
232 continue;
233 }
234
235 if matches!(state_update.kind, StateUpdateKind::SystemConfiguration(_)) {
236 update_system_config = true;
237 }
238
239 match state_update.diff {
240 StateDiff::Retraction => {
241 if let Some(update) =
245 parsed_state_updates::parse_state_update(self, state_update.clone())
246 {
247 catalog_updates.push(update);
248 }
249
250 builtin_table_updates.extend(self.generate_builtin_table_update(
253 state_update.kind.clone(),
254 state_update.diff,
255 ));
256 self.apply_update(
257 state_update.kind,
258 state_update.diff,
259 retractions,
260 local_expression_cache,
261 )?;
262 }
263 StateDiff::Addition => {
264 self.apply_update(
265 state_update.kind.clone(),
266 state_update.diff,
267 retractions,
268 local_expression_cache,
269 )?;
270 builtin_table_updates.extend(self.generate_builtin_table_update(
274 state_update.kind.clone(),
275 state_update.diff,
276 ));
277
278 if let Some(update) =
281 parsed_state_updates::parse_state_update(self, state_update.clone())
282 {
283 catalog_updates.push(update);
284 }
285 }
286 }
287 }
288
289 if update_system_config {
290 self.system_configuration.sync_dyncfgs();
291 }
292
293 Ok((builtin_table_updates, catalog_updates))
294 }
295
296 #[instrument(level = "debug")]
297 fn apply_update(
298 &mut self,
299 kind: StateUpdateKind,
300 diff: StateDiff,
301 retractions: &mut InProgressRetractions,
302 local_expression_cache: &mut LocalExpressionCache,
303 ) -> Result<(), CatalogError> {
304 match kind {
305 StateUpdateKind::Role(role) => {
306 self.apply_role_update(role, diff, retractions);
307 }
308 StateUpdateKind::RoleAuth(role_auth) => {
309 self.apply_role_auth_update(role_auth, diff, retractions);
310 }
311 StateUpdateKind::Database(database) => {
312 self.apply_database_update(database, diff, retractions);
313 }
314 StateUpdateKind::Schema(schema) => {
315 self.apply_schema_update(schema, diff, retractions);
316 }
317 StateUpdateKind::DefaultPrivilege(default_privilege) => {
318 self.apply_default_privilege_update(default_privilege, diff, retractions);
319 }
320 StateUpdateKind::SystemPrivilege(system_privilege) => {
321 self.apply_system_privilege_update(system_privilege, diff, retractions);
322 }
323 StateUpdateKind::SystemConfiguration(system_configuration) => {
324 self.apply_system_configuration_update(system_configuration, diff, retractions);
325 }
326 StateUpdateKind::ClusterSystemConfiguration(cfg) => {
327 Self::apply_scoped_system_configuration_update(
328 &mut self.scoped_system_parameters.cluster,
329 cfg.cluster_id,
330 cfg.name,
331 cfg.value,
332 diff,
333 );
334 }
335 StateUpdateKind::ReplicaSystemConfiguration(cfg) => {
336 Self::apply_scoped_system_configuration_update(
337 &mut self.scoped_system_parameters.replica,
338 cfg.replica_id,
339 cfg.name,
340 cfg.value,
341 diff,
342 );
343 }
344 StateUpdateKind::Cluster(cluster) => {
345 self.apply_cluster_update(cluster, diff, retractions);
346 }
347 StateUpdateKind::NetworkPolicy(network_policy) => {
348 self.apply_network_policy_update(network_policy, diff, retractions);
349 }
350 StateUpdateKind::IntrospectionSourceIndex(introspection_source_index) => {
351 self.apply_introspection_source_index_update(
352 introspection_source_index,
353 diff,
354 retractions,
355 );
356 }
357 StateUpdateKind::ClusterReplica(cluster_replica) => {
358 self.apply_cluster_replica_update(cluster_replica, diff, retractions);
359 }
360 StateUpdateKind::SystemObjectMapping(system_object_mapping) => {
361 self.apply_system_object_mapping_update(
362 system_object_mapping,
363 diff,
364 retractions,
365 local_expression_cache,
366 );
367 }
368 StateUpdateKind::Item(item) => {
369 self.apply_item_update(item, diff, retractions, local_expression_cache)?;
370 }
371 StateUpdateKind::Comment(comment) => {
372 self.apply_comment_update(comment, diff, retractions);
373 }
374 StateUpdateKind::SourceReferences(source_reference) => {
375 self.apply_source_references_update(source_reference, diff, retractions);
376 }
377 StateUpdateKind::AuditLog(_audit_log) => {
378 }
380 StateUpdateKind::StorageCollectionMetadata(storage_collection_metadata) => {
381 self.apply_storage_collection_metadata_update(
382 storage_collection_metadata,
383 diff,
384 retractions,
385 );
386 }
387 StateUpdateKind::UnfinalizedShard(unfinalized_shard) => {
388 self.apply_unfinalized_shard_update(unfinalized_shard, diff, retractions);
389 }
390 }
391
392 Ok(())
393 }
394
395 #[instrument(level = "debug")]
396 fn apply_role_auth_update(
397 &mut self,
398 role_auth: mz_catalog::durable::RoleAuth,
399 diff: StateDiff,
400 retractions: &mut InProgressRetractions,
401 ) {
402 apply_with_update(
403 &mut self.role_auth_by_id,
404 role_auth,
405 |role_auth| role_auth.role_id,
406 diff,
407 &mut retractions.role_auths,
408 );
409 }
410
411 #[instrument(level = "debug")]
412 fn apply_role_update(
413 &mut self,
414 role: mz_catalog::durable::Role,
415 diff: StateDiff,
416 retractions: &mut InProgressRetractions,
417 ) {
418 apply_inverted_lookup(&mut self.roles_by_name, &role.name, role.id, diff);
419 apply_with_update(
420 &mut self.roles_by_id,
421 role,
422 |role| role.id,
423 diff,
424 &mut retractions.roles,
425 );
426 }
427
428 #[instrument(level = "debug")]
429 fn apply_database_update(
430 &mut self,
431 database: mz_catalog::durable::Database,
432 diff: StateDiff,
433 retractions: &mut InProgressRetractions,
434 ) {
435 apply_inverted_lookup(
436 &mut self.database_by_name,
437 &database.name,
438 database.id,
439 diff,
440 );
441 apply_with_update(
442 &mut self.database_by_id,
443 database,
444 |database| database.id,
445 diff,
446 &mut retractions.databases,
447 );
448 }
449
450 #[instrument(level = "debug")]
451 fn apply_schema_update(
452 &mut self,
453 schema: mz_catalog::durable::Schema,
454 diff: StateDiff,
455 retractions: &mut InProgressRetractions,
456 ) {
457 match &schema.database_id {
458 Some(database_id) => {
459 let db = self
460 .database_by_id
461 .get_mut(database_id)
462 .expect("catalog out of sync");
463 apply_inverted_lookup(&mut db.schemas_by_name, &schema.name, schema.id, diff);
464 apply_with_update(
465 &mut db.schemas_by_id,
466 schema,
467 |schema| schema.id,
468 diff,
469 &mut retractions.schemas,
470 );
471 }
472 None => {
473 apply_inverted_lookup(
474 &mut self.ambient_schemas_by_name,
475 &schema.name,
476 schema.id,
477 diff,
478 );
479 apply_with_update(
480 &mut self.ambient_schemas_by_id,
481 schema,
482 |schema| schema.id,
483 diff,
484 &mut retractions.schemas,
485 );
486 }
487 }
488 }
489
490 #[instrument(level = "debug")]
491 fn apply_default_privilege_update(
492 &mut self,
493 default_privilege: mz_catalog::durable::DefaultPrivilege,
494 diff: StateDiff,
495 _retractions: &mut InProgressRetractions,
496 ) {
497 match diff {
498 StateDiff::Addition => Arc::make_mut(&mut self.default_privileges)
499 .grant(default_privilege.object, default_privilege.acl_item),
500 StateDiff::Retraction => Arc::make_mut(&mut self.default_privileges)
501 .revoke(&default_privilege.object, &default_privilege.acl_item),
502 }
503 }
504
505 #[instrument(level = "debug")]
506 fn apply_system_privilege_update(
507 &mut self,
508 system_privilege: MzAclItem,
509 diff: StateDiff,
510 _retractions: &mut InProgressRetractions,
511 ) {
512 match diff {
513 StateDiff::Addition => {
514 Arc::make_mut(&mut self.system_privileges).grant(system_privilege)
515 }
516 StateDiff::Retraction => {
517 Arc::make_mut(&mut self.system_privileges).revoke(&system_privilege)
518 }
519 }
520 }
521
522 #[instrument(level = "debug")]
523 fn apply_system_configuration_update(
524 &mut self,
525 system_configuration: mz_catalog::durable::SystemConfiguration,
526 diff: StateDiff,
527 _retractions: &mut InProgressRetractions,
528 ) {
529 let res = match diff {
530 StateDiff::Addition => self.insert_system_configuration(
531 &system_configuration.name,
532 VarInput::Flat(&system_configuration.value),
533 ),
534 StateDiff::Retraction => self.remove_system_configuration(&system_configuration.name),
535 };
536 match res {
537 Ok(_) => (),
538 Err(Error {
542 kind: ErrorKind::VarError(VarError::UnknownParameter(name)),
543 }) => {
544 warn!(%name, "unknown system parameter from catalog storage");
545 }
546 Err(e) => panic!("unable to update system variable: {e:?}"),
547 }
548 }
549
550 fn apply_scoped_system_configuration_update<Id: Ord>(
564 map: &mut BTreeMap<Id, BTreeMap<String, String>>,
565 id: Id,
566 name: String,
567 value: String,
568 diff: StateDiff,
569 ) {
570 match diff {
571 StateDiff::Addition => {
572 map.entry(id).or_default().insert(name, value);
573 }
574 StateDiff::Retraction => {
575 if let Some(values) = map.get_mut(&id) {
576 if values.get(&name) == Some(&value) {
577 values.remove(&name);
578 if values.is_empty() {
579 map.remove(&id);
580 }
581 }
582 }
583 }
584 }
585 }
586
587 #[instrument(level = "debug")]
588 fn apply_cluster_update(
589 &mut self,
590 cluster: mz_catalog::durable::Cluster,
591 diff: StateDiff,
592 retractions: &mut InProgressRetractions,
593 ) {
594 if matches!(diff, StateDiff::Addition) {
601 if let mz_catalog::durable::ClusterVariant::Managed(managed) = &cluster.config.variant {
602 if !self.cluster_replica_sizes.0.contains_key(&managed.size) {
603 soft_panic_or_log!(
604 "managed cluster {} ({}) references unknown replica size {:?}; \
605 mz_clusters.disk will resolve to false",
606 cluster.name,
607 cluster.id,
608 managed.size,
609 );
610 }
611 }
612 }
613 apply_inverted_lookup(&mut self.clusters_by_name, &cluster.name, cluster.id, diff);
614 apply_with_update(
615 &mut self.clusters_by_id,
616 cluster,
617 |cluster| cluster.id,
618 diff,
619 &mut retractions.clusters,
620 );
621 }
622
623 #[instrument(level = "debug")]
624 fn apply_network_policy_update(
625 &mut self,
626 policy: mz_catalog::durable::NetworkPolicy,
627 diff: StateDiff,
628 retractions: &mut InProgressRetractions,
629 ) {
630 apply_inverted_lookup(
631 &mut self.network_policies_by_name,
632 &policy.name,
633 policy.id,
634 diff,
635 );
636 apply_with_update(
637 &mut self.network_policies_by_id,
638 policy,
639 |policy| policy.id,
640 diff,
641 &mut retractions.network_policies,
642 );
643 }
644
645 #[instrument(level = "debug")]
646 fn apply_introspection_source_index_update(
647 &mut self,
648 introspection_source_index: mz_catalog::durable::IntrospectionSourceIndex,
649 diff: StateDiff,
650 retractions: &mut InProgressRetractions,
651 ) {
652 let cluster = self
653 .clusters_by_id
654 .get_mut(&introspection_source_index.cluster_id)
655 .expect("catalog out of sync");
656 let log = BUILTIN_LOG_LOOKUP
657 .get(introspection_source_index.name.as_str())
658 .expect("missing log");
659 apply_inverted_lookup(
660 &mut cluster.log_indexes,
661 &log.variant,
662 introspection_source_index.index_id,
663 diff,
664 );
665
666 match diff {
667 StateDiff::Addition => {
668 if let Some(mut entry) = retractions
669 .introspection_source_indexes
670 .remove(&introspection_source_index.item_id)
671 {
672 let (index_name, index) = self.create_introspection_source_index(
675 introspection_source_index.cluster_id,
676 log,
677 introspection_source_index.index_id,
678 );
679 assert_eq!(entry.id, introspection_source_index.item_id);
680 assert_eq!(entry.oid, introspection_source_index.oid);
681 assert_eq!(entry.name, index_name);
682 entry.item = index;
683 self.insert_entry(entry);
684 } else {
685 self.insert_introspection_source_index(
686 introspection_source_index.cluster_id,
687 log,
688 introspection_source_index.item_id,
689 introspection_source_index.index_id,
690 introspection_source_index.oid,
691 );
692 }
693 }
694 StateDiff::Retraction => {
695 let entry = self.drop_item(introspection_source_index.item_id);
696 retractions
697 .introspection_source_indexes
698 .insert(entry.id, entry);
699 }
700 }
701 }
702
703 #[instrument(level = "debug")]
704 fn apply_cluster_replica_update(
705 &mut self,
706 cluster_replica: mz_catalog::durable::ClusterReplica,
707 diff: StateDiff,
708 _retractions: &mut InProgressRetractions,
709 ) {
710 let cluster = self
711 .clusters_by_id
712 .get(&cluster_replica.cluster_id)
713 .expect("catalog out of sync");
714
715 if let mz_catalog::durable::ReplicaLocation::Managed { size, .. } =
724 &cluster_replica.config.location
725 {
726 if !self.cluster_replica_sizes.0.contains_key(size) {
727 soft_panic_or_log!(
728 "cluster replica {}.{} ({}) references unknown replica size {:?}; \
729 skipping in-memory registration",
730 cluster.name,
731 cluster_replica.name,
732 cluster_replica.replica_id,
733 size,
734 );
735 return;
736 }
737 }
738
739 let location = self
745 .concretize_replica_location(cluster_replica.config.location, &vec![], None, true)
746 .expect("catalog in unexpected state");
747 let cluster = self
748 .clusters_by_id
749 .get_mut(&cluster_replica.cluster_id)
750 .expect("catalog out of sync");
751 apply_inverted_lookup(
752 &mut cluster.replica_id_by_name_,
753 &cluster_replica.name,
754 cluster_replica.replica_id,
755 diff,
756 );
757 match diff {
758 StateDiff::Retraction => {
759 let prev = cluster.replicas_by_id_.remove(&cluster_replica.replica_id);
760 assert!(
761 prev.is_some(),
762 "retraction does not match existing value: {:?}",
763 cluster_replica.replica_id
764 );
765 }
766 StateDiff::Addition => {
767 let logging = ReplicaLogging {
768 log_logging: cluster_replica.config.logging.log_logging,
769 interval: cluster_replica.config.logging.interval,
770 };
771 let config = ReplicaConfig {
772 location,
773 compute: ComputeReplicaConfig {
774 logging,
775 arrangement_compression: cluster_replica.config.arrangement_compression,
776 },
777 };
778 let mem_cluster_replica = ClusterReplica {
779 name: cluster_replica.name.clone(),
780 cluster_id: cluster_replica.cluster_id,
781 replica_id: cluster_replica.replica_id,
782 config,
783 owner_id: cluster_replica.owner_id,
784 };
785 let prev = cluster
786 .replicas_by_id_
787 .insert(cluster_replica.replica_id, mem_cluster_replica);
788 assert_eq!(
789 prev, None,
790 "values must be explicitly retracted before inserting a new value: {:?}",
791 cluster_replica.replica_id
792 );
793 }
794 }
795 }
796
797 #[instrument(level = "debug")]
798 fn apply_system_object_mapping_update(
799 &mut self,
800 system_object_mapping: mz_catalog::durable::SystemObjectMapping,
801 diff: StateDiff,
802 retractions: &mut InProgressRetractions,
803 local_expression_cache: &mut LocalExpressionCache,
804 ) {
805 let item_id = system_object_mapping.unique_identifier.catalog_id;
806 let global_id = system_object_mapping.unique_identifier.global_id;
807
808 if system_object_mapping.unique_identifier.runtime_alterable() {
809 return;
813 }
814
815 if let StateDiff::Retraction = diff {
816 let entry = self.drop_item(item_id);
817 retractions.system_object_mappings.insert(item_id, entry);
818 return;
819 }
820
821 if let Some(entry) = retractions.system_object_mappings.remove(&item_id) {
822 self.insert_entry(entry);
827 return;
828 }
829
830 let builtin = BUILTIN_LOOKUP
831 .get(&system_object_mapping.description)
832 .expect("missing builtin")
833 .1;
834 let schema_name = builtin.schema();
835 let schema_id = self
836 .ambient_schemas_by_name
837 .get(schema_name)
838 .unwrap_or_else(|| panic!("unknown ambient schema: {schema_name}"));
839 let name = QualifiedItemName {
840 qualifiers: ItemQualifiers {
841 database_spec: ResolvedDatabaseSpecifier::Ambient,
842 schema_spec: SchemaSpecifier::Id(*schema_id),
843 },
844 item: builtin.name().into(),
845 };
846 match builtin {
847 Builtin::Log(log) => {
848 let mut acl_items = vec![rbac::owner_privilege(
849 mz_sql::catalog::ObjectType::Source,
850 MZ_SYSTEM_ROLE_ID,
851 )];
852 acl_items.extend_from_slice(&log.access);
853 self.insert_item(
854 item_id,
855 log.oid,
856 name.clone(),
857 CatalogItem::Log(Log {
858 variant: log.variant,
859 global_id,
860 }),
861 MZ_SYSTEM_ROLE_ID,
862 PrivilegeMap::from_mz_acl_items(acl_items),
863 );
864 }
865
866 Builtin::Table(table) => {
867 let mut acl_items = vec![rbac::owner_privilege(
868 mz_sql::catalog::ObjectType::Table,
869 MZ_SYSTEM_ROLE_ID,
870 )];
871 acl_items.extend_from_slice(&table.access);
872
873 self.insert_item(
874 item_id,
875 table.oid,
876 name.clone(),
877 CatalogItem::Table(Table {
878 create_sql: None,
879 desc: VersionedRelationDesc::new(table.desc.clone()),
880 collections: [(RelationVersion::root(), global_id)].into_iter().collect(),
881 conn_id: None,
882 resolved_ids: ResolvedIds::empty(),
883 custom_logical_compaction_window: table.is_retained_metrics_object.then(
884 || {
885 self.system_config()
886 .metrics_retention()
887 .try_into()
888 .expect("invalid metrics retention")
889 },
890 ),
891 is_retained_metrics_object: table.is_retained_metrics_object,
892 data_source: TableDataSource::TableWrites {
893 defaults: vec![Expr::null(); table.desc.arity()],
894 },
895 }),
896 MZ_SYSTEM_ROLE_ID,
897 PrivilegeMap::from_mz_acl_items(acl_items),
898 );
899 }
900 Builtin::Index(index) => {
901 let custom_logical_compaction_window =
902 index.is_retained_metrics_object.then(|| {
903 self.system_config()
904 .metrics_retention()
905 .try_into()
906 .expect("invalid metrics retention")
907 });
908 let versions = BTreeMap::new();
910
911 let item = self
912 .parse_item(
913 global_id,
914 &index.create_sql(),
915 &versions,
916 None,
917 index.is_retained_metrics_object,
918 custom_logical_compaction_window,
919 local_expression_cache,
920 None,
921 )
922 .unwrap_or_else(|e| {
923 panic!(
924 "internal error: failed to load bootstrap index:\n\
925 {}\n\
926 error:\n\
927 {:?}\n\n\
928 make sure that the schema name is specified in the builtin index's create sql statement.",
929 index.name, e
930 )
931 });
932 let CatalogItem::Index(_) = item else {
933 panic!(
934 "internal error: builtin index {}'s SQL does not begin with \"CREATE INDEX\".",
935 index.name
936 );
937 };
938
939 self.insert_item(
940 item_id,
941 index.oid,
942 name,
943 item,
944 MZ_SYSTEM_ROLE_ID,
945 PrivilegeMap::default(),
946 );
947 }
948 Builtin::View(_) => {
949 unreachable!("views added elsewhere");
951 }
952
953 Builtin::Type(typ) => {
955 let typ = self.resolve_builtin_type_references(typ);
956 if let CatalogType::Array { element_reference } = typ.details.typ {
957 let entry = self.get_entry_mut(&element_reference);
958 let item_type = match &mut entry.item {
959 CatalogItem::Type(item_type) => item_type,
960 _ => unreachable!("types can only reference other types"),
961 };
962 item_type.details.array_id = Some(item_id);
963 }
964
965 let schema_id = self.resolve_system_schema(typ.schema);
966
967 self.insert_item(
968 item_id,
969 typ.oid,
970 QualifiedItemName {
971 qualifiers: ItemQualifiers {
972 database_spec: ResolvedDatabaseSpecifier::Ambient,
973 schema_spec: SchemaSpecifier::Id(schema_id),
974 },
975 item: typ.name.to_owned(),
976 },
977 CatalogItem::Type(Type {
978 create_sql: None,
979 global_id,
980 details: typ.details.clone(),
981 resolved_ids: ResolvedIds::empty(),
982 }),
983 MZ_SYSTEM_ROLE_ID,
984 PrivilegeMap::from_mz_acl_items(vec![
985 rbac::default_builtin_object_privilege(mz_sql::catalog::ObjectType::Type),
986 rbac::owner_privilege(mz_sql::catalog::ObjectType::Type, MZ_SYSTEM_ROLE_ID),
987 ]),
988 );
989 }
990
991 Builtin::Func(func) => {
992 let oid = INVALID_OID;
996 self.insert_item(
997 item_id,
998 oid,
999 name.clone(),
1000 CatalogItem::Func(Func {
1001 inner: func.inner,
1002 global_id,
1003 }),
1004 MZ_SYSTEM_ROLE_ID,
1005 PrivilegeMap::default(),
1006 );
1007 }
1008
1009 Builtin::Source(coll) => {
1010 let mut acl_items = vec![rbac::owner_privilege(
1011 mz_sql::catalog::ObjectType::Source,
1012 MZ_SYSTEM_ROLE_ID,
1013 )];
1014 acl_items.extend_from_slice(&coll.access);
1015
1016 self.insert_item(
1017 item_id,
1018 coll.oid,
1019 name.clone(),
1020 CatalogItem::Source(Source {
1021 create_sql: None,
1022 data_source: coll.data_source.clone(),
1023 desc: coll.desc.clone(),
1024 global_id,
1025 timeline: Timeline::EpochMilliseconds,
1026 resolved_ids: ResolvedIds::empty(),
1027 custom_logical_compaction_window: coll.is_retained_metrics_object.then(
1028 || {
1029 self.system_config()
1030 .metrics_retention()
1031 .try_into()
1032 .expect("invalid metrics retention")
1033 },
1034 ),
1035 is_retained_metrics_object: coll.is_retained_metrics_object,
1036 }),
1037 MZ_SYSTEM_ROLE_ID,
1038 PrivilegeMap::from_mz_acl_items(acl_items),
1039 );
1040 }
1041 Builtin::MaterializedView(mv) => {
1042 let mut acl_items = vec![rbac::owner_privilege(
1043 mz_sql::catalog::ObjectType::MaterializedView,
1044 MZ_SYSTEM_ROLE_ID,
1045 )];
1046 acl_items.extend_from_slice(&mv.access);
1047
1048 let custom_logical_compaction_window = mv.is_retained_metrics_object.then(|| {
1049 self.system_config()
1050 .metrics_retention()
1051 .try_into()
1052 .expect("invalid metrics retention")
1053 });
1054
1055 let versions = BTreeMap::new();
1057
1058 let mut item = self
1059 .parse_item(
1060 global_id,
1061 &mv.create_sql(),
1062 &versions,
1063 None,
1064 mv.is_retained_metrics_object,
1065 custom_logical_compaction_window,
1066 local_expression_cache,
1067 None,
1068 )
1069 .unwrap_or_else(|e| {
1070 panic!(
1071 "internal error: failed to load bootstrap materialized view:\n\
1072 {}\n\
1073 error:\n\
1074 {e:?}\n\n\
1075 make sure that the schema name is specified in the builtin \
1076 materialized view's create sql statement.",
1077 mv.name,
1078 )
1079 });
1080 let CatalogItem::MaterializedView(catalog_mv) = &mut item else {
1081 panic!(
1082 "internal error: builtin materialized view {}'s SQL does not begin \
1083 with \"CREATE MATERIALIZED VIEW\".",
1084 mv.name,
1085 );
1086 };
1087
1088 let mut desc = catalog_mv.desc.latest();
1092 for key in &mv.desc.typ().keys {
1093 desc = desc.with_key(key.clone());
1094 }
1095 catalog_mv.desc = VersionedRelationDesc::new(desc);
1096
1097 self.insert_item(
1098 item_id,
1099 mv.oid,
1100 name,
1101 item,
1102 MZ_SYSTEM_ROLE_ID,
1103 PrivilegeMap::from_mz_acl_items(acl_items),
1104 );
1105 }
1106 Builtin::Connection(connection) => {
1107 let versions = BTreeMap::new();
1109 let mut item = self
1110 .parse_item(
1111 global_id,
1112 connection.sql,
1113 &versions,
1114 None,
1115 false,
1116 None,
1117 local_expression_cache,
1118 None,
1119 )
1120 .unwrap_or_else(|e| {
1121 panic!(
1122 "internal error: failed to load bootstrap connection:\n\
1123 {}\n\
1124 error:\n\
1125 {:?}\n\n\
1126 make sure that the schema name is specified in the builtin connection's create sql statement.",
1127 connection.name, e
1128 )
1129 });
1130 let CatalogItem::Connection(_) = &mut item else {
1131 panic!(
1132 "internal error: builtin connection {}'s SQL does not begin with \"CREATE CONNECTION\".",
1133 connection.name
1134 );
1135 };
1136
1137 let mut acl_items = vec![rbac::owner_privilege(
1138 mz_sql::catalog::ObjectType::Connection,
1139 connection.owner_id.clone(),
1140 )];
1141 acl_items.extend_from_slice(connection.access);
1142
1143 self.insert_item(
1144 item_id,
1145 connection.oid,
1146 name.clone(),
1147 item,
1148 connection.owner_id.clone(),
1149 PrivilegeMap::from_mz_acl_items(acl_items),
1150 );
1151 }
1152 }
1153 }
1154
1155 fn is_nonlocal_ephemeral_item_update(&self, update: &StateUpdate) -> bool {
1158 let StateUpdateKind::Item(item) = &update.kind else {
1159 return false;
1160 };
1161 let Some(owner) = item.ephemeral_owner_session else {
1163 return false;
1164 };
1165 match update.diff {
1166 StateDiff::Addition => !self.temporary_namespaces.contains_uuid(&owner),
1168 StateDiff::Retraction => !self.entry_by_id.contains_key(&item.id),
1171 }
1172 }
1173
1174 #[instrument(level = "debug")]
1175 fn apply_item_update(
1176 &mut self,
1177 item: mz_catalog::durable::Item,
1178 diff: StateDiff,
1179 retractions: &mut InProgressRetractions,
1180 local_expression_cache: &mut LocalExpressionCache,
1181 ) -> Result<(), CatalogError> {
1182 match diff {
1183 StateDiff::Addition => {
1184 let key = item.key();
1185 let mz_catalog::durable::Item {
1186 id,
1187 oid,
1188 global_id,
1189 schema_id,
1190 name,
1191 create_sql,
1192 owner_id,
1193 privileges,
1194 extra_versions,
1195 ephemeral_owner_session,
1196 } = item;
1197
1198 let conn_id = ephemeral_owner_session.map(|owner| {
1205 self.temporary_namespaces
1206 .conn_for_uuid(&owner)
1207 .cloned()
1208 .unwrap_or_else(|| {
1209 panic!("no session record applied for temporary item owner {owner}")
1210 })
1211 });
1212 let name = match &conn_id {
1213 Some(_) => QualifiedItemName {
1214 qualifiers: ItemQualifiers {
1215 database_spec: ResolvedDatabaseSpecifier::Ambient,
1216 schema_spec: SchemaSpecifier::Temporary,
1217 },
1218 item: name.clone(),
1219 },
1220 None => {
1221 let schema = self.find_non_temp_schema(&schema_id);
1222 QualifiedItemName {
1223 qualifiers: ItemQualifiers {
1224 database_spec: schema.database().clone(),
1225 schema_spec: schema.id().clone(),
1226 },
1227 item: name.clone(),
1228 }
1229 }
1230 };
1231 let entry = match retractions.items.remove(&key) {
1232 Some(mut retraction) => {
1233 assert_eq!(retraction.id, id);
1234
1235 if retraction.create_sql() != create_sql {
1240 let mut catalog_item = self
1241 .deserialize_item(
1242 global_id,
1243 &create_sql,
1244 &extra_versions,
1245 local_expression_cache,
1246 Some(retraction.item),
1247 )
1248 .unwrap_or_else(|e| {
1249 panic!("{e:?}: invalid persisted SQL: {create_sql}")
1250 });
1251 if conn_id.is_some() {
1252 catalog_item.set_conn_id(conn_id.clone());
1260 catalog_item.set_create_sql(create_sql);
1272 }
1273 retraction.item = catalog_item;
1274 }
1275
1276 retraction.id = id;
1277 retraction.oid = oid;
1278 retraction.name = name;
1279 retraction.owner_id = owner_id;
1280 retraction.privileges = PrivilegeMap::from_mz_acl_items(privileges);
1281 retraction
1282 }
1283 None => {
1284 let mut catalog_item = self
1285 .deserialize_item(
1286 global_id,
1287 &create_sql,
1288 &extra_versions,
1289 local_expression_cache,
1290 None,
1291 )
1292 .unwrap_or_else(|e| {
1293 panic!("{e:?}: invalid persisted SQL: {create_sql}")
1294 });
1295
1296 if conn_id.is_some() {
1297 catalog_item.set_conn_id(conn_id.clone());
1299 catalog_item.set_create_sql(create_sql);
1300 }
1301
1302 CatalogEntry {
1303 item: catalog_item,
1304 referenced_by: Vec::new(),
1305 used_by: Vec::new(),
1306 id,
1307 oid,
1308 name,
1309 owner_id,
1310 privileges: PrivilegeMap::from_mz_acl_items(privileges),
1311 }
1312 }
1313 };
1314
1315 self.insert_entry(entry);
1316 }
1317 StateDiff::Retraction => {
1318 let entry = self.drop_item(item.id);
1319 let key = item.into_key_value().0;
1320 retractions.items.insert(key, entry);
1321 }
1322 }
1323 Ok(())
1324 }
1325
1326 #[instrument(level = "debug")]
1327 fn apply_comment_update(
1328 &mut self,
1329 comment: mz_catalog::durable::Comment,
1330 diff: StateDiff,
1331 _retractions: &mut InProgressRetractions,
1332 ) {
1333 match diff {
1334 StateDiff::Addition => {
1335 let prev = Arc::make_mut(&mut self.comments).update_comment(
1336 comment.object_id,
1337 comment.sub_component,
1338 Some(comment.comment),
1339 );
1340 assert_eq!(
1341 prev, None,
1342 "values must be explicitly retracted before inserting a new value"
1343 );
1344 }
1345 StateDiff::Retraction => {
1346 let prev = Arc::make_mut(&mut self.comments).update_comment(
1347 comment.object_id,
1348 comment.sub_component,
1349 None,
1350 );
1351 assert_eq!(
1352 prev,
1353 Some(comment.comment),
1354 "retraction does not match existing value: ({:?}, {:?})",
1355 comment.object_id,
1356 comment.sub_component,
1357 );
1358 }
1359 }
1360 }
1361
1362 #[instrument(level = "debug")]
1363 fn apply_source_references_update(
1364 &mut self,
1365 source_references: mz_catalog::durable::SourceReferences,
1366 diff: StateDiff,
1367 _retractions: &mut InProgressRetractions,
1368 ) {
1369 match diff {
1370 StateDiff::Addition => {
1371 let prev = self
1372 .source_references
1373 .insert(source_references.source_id, source_references.into());
1374 assert!(
1375 prev.is_none(),
1376 "values must be explicitly retracted before inserting a new value: {prev:?}"
1377 );
1378 }
1379 StateDiff::Retraction => {
1380 let prev = self.source_references.remove(&source_references.source_id);
1381 assert!(
1382 prev.is_some(),
1383 "retraction for a non-existent existing value: {source_references:?}"
1384 );
1385 }
1386 }
1387 }
1388
1389 #[instrument(level = "debug")]
1390 fn apply_storage_collection_metadata_update(
1391 &mut self,
1392 storage_collection_metadata: mz_catalog::durable::StorageCollectionMetadata,
1393 diff: StateDiff,
1394 _retractions: &mut InProgressRetractions,
1395 ) {
1396 apply_inverted_lookup(
1397 &mut Arc::make_mut(&mut self.storage_metadata).collection_metadata,
1398 &storage_collection_metadata.id,
1399 storage_collection_metadata.shard,
1400 diff,
1401 );
1402 }
1403
1404 #[instrument(level = "debug")]
1405 fn apply_unfinalized_shard_update(
1406 &mut self,
1407 unfinalized_shard: mz_catalog::durable::UnfinalizedShard,
1408 diff: StateDiff,
1409 _retractions: &mut InProgressRetractions,
1410 ) {
1411 match diff {
1412 StateDiff::Addition => {
1413 let newly_inserted = Arc::make_mut(&mut self.storage_metadata)
1414 .unfinalized_shards
1415 .insert(unfinalized_shard.shard);
1416 assert!(
1417 newly_inserted,
1418 "values must be explicitly retracted before inserting a new value: {unfinalized_shard:?}",
1419 );
1420 }
1421 StateDiff::Retraction => {
1422 let removed = Arc::make_mut(&mut self.storage_metadata)
1423 .unfinalized_shards
1424 .remove(&unfinalized_shard.shard);
1425 assert!(
1426 removed,
1427 "retraction does not match existing value: {unfinalized_shard:?}"
1428 );
1429 }
1430 }
1431 }
1432
1433 #[instrument(level = "debug")]
1436 pub(crate) fn generate_builtin_table_update(
1437 &self,
1438 kind: StateUpdateKind,
1439 diff: StateDiff,
1440 ) -> Vec<BuiltinTableUpdate<&'static BuiltinTable>> {
1441 let diff = diff.into();
1442 match kind {
1443 StateUpdateKind::Role(_) => Vec::new(),
1447 StateUpdateKind::RoleAuth(role_auth) => {
1448 vec![self.pack_role_auth_update(role_auth.role_id, diff)]
1449 }
1450 StateUpdateKind::DefaultPrivilege(_) => Vec::new(),
1454 StateUpdateKind::SystemPrivilege(_) => Vec::new(),
1455 StateUpdateKind::SystemConfiguration(_) => Vec::new(),
1456 StateUpdateKind::ClusterSystemConfiguration(_) => Vec::new(),
1462 StateUpdateKind::ReplicaSystemConfiguration(_) => Vec::new(),
1463 StateUpdateKind::Cluster(_) => Vec::new(),
1467 StateUpdateKind::IntrospectionSourceIndex(introspection_source_index) => {
1468 self.pack_item_update(introspection_source_index.item_id, diff)
1469 }
1470 StateUpdateKind::ClusterReplica(_) => Vec::new(),
1473 StateUpdateKind::SystemObjectMapping(system_object_mapping) => {
1474 if !system_object_mapping.unique_identifier.runtime_alterable() {
1478 self.pack_item_update(system_object_mapping.unique_identifier.catalog_id, diff)
1479 } else {
1480 vec![]
1481 }
1482 }
1483 StateUpdateKind::Item(item) => self.pack_item_update(item.id, diff),
1484 StateUpdateKind::Comment(_) => Vec::new(),
1485 StateUpdateKind::SourceReferences(source_references) => {
1486 self.pack_source_references_update(&source_references, diff)
1487 }
1488 StateUpdateKind::AuditLog(_) => Vec::new(),
1492 StateUpdateKind::Database(_)
1493 | StateUpdateKind::Schema(_)
1494 | StateUpdateKind::NetworkPolicy(_)
1495 | StateUpdateKind::StorageCollectionMetadata(_)
1496 | StateUpdateKind::UnfinalizedShard(_) => Vec::new(),
1497 }
1498 }
1499
1500 fn get_entry_mut(&mut self, id: &CatalogItemId) -> &mut CatalogEntry {
1501 self.entry_by_id
1502 .get_mut(id)
1503 .unwrap_or_else(|| panic!("catalog out of sync, missing id {id}"))
1504 }
1505
1506 pub(super) fn set_optimized_plan(
1511 &mut self,
1512 id: GlobalId,
1513 plan: DataflowDescription<mz_expr::OptimizedMirRelationExpr>,
1514 ) {
1515 let item_id = self.entry_by_global_id[&id];
1516 let entry = self.get_entry_mut(&item_id);
1517 match entry.item_mut() {
1518 CatalogItem::Index(idx) => idx.optimized_plan = Some(Arc::new(plan)),
1519 CatalogItem::MaterializedView(mv) => mv.optimized_plan = Some(Arc::new(plan)),
1520 CatalogItem::MetricSink(ms) => ms.optimized_plan = Some(Arc::new(plan)),
1521 other => panic!("set_optimized_plan called on {} ({:?})", id, other.typ()),
1522 }
1523 }
1524
1525 pub(super) fn set_physical_plan(
1530 &mut self,
1531 id: GlobalId,
1532 plan: DataflowDescription<mz_compute_types::plan::LirRelationExpr>,
1533 ) {
1534 let item_id = self.entry_by_global_id[&id];
1535 let entry = self.get_entry_mut(&item_id);
1536 match entry.item_mut() {
1537 CatalogItem::Index(idx) => idx.physical_plan = Some(Arc::new(plan)),
1538 CatalogItem::MaterializedView(mv) => mv.physical_plan = Some(Arc::new(plan)),
1539 CatalogItem::MetricSink(ms) => ms.physical_plan = Some(Arc::new(plan)),
1540 other => panic!("set_physical_plan called on {} ({:?})", id, other.typ()),
1541 }
1542 }
1543
1544 pub(super) fn set_dataflow_metainfo(
1549 &mut self,
1550 id: GlobalId,
1551 metainfo: DataflowMetainfo<Arc<OptimizerNotice>>,
1552 ) {
1553 for notice in metainfo.optimizer_notices.iter() {
1555 for dep_id in notice.dependencies.iter() {
1556 self.notices_by_dep_id
1557 .entry(*dep_id)
1558 .or_default()
1559 .push(Arc::clone(notice));
1560 }
1561 if let Some(item_id) = notice.item_id {
1562 soft_assert_eq_or_log!(
1563 item_id,
1564 id,
1565 "notice.item_id should match the id for whom we are saving the notice"
1566 );
1567 }
1568 }
1569 let item_id = self.entry_by_global_id[&id];
1571 let entry = self.get_entry_mut(&item_id);
1572 match entry.item_mut() {
1573 CatalogItem::Index(idx) => idx.dataflow_metainfo = Some(metainfo),
1574 CatalogItem::MaterializedView(mv) => mv.dataflow_metainfo = Some(metainfo),
1575 CatalogItem::MetricSink(ms) => ms.dataflow_metainfo = Some(metainfo),
1576 other => panic!("set_dataflow_metainfo called on {} ({:?})", id, other.typ()),
1577 }
1578 }
1579
1580 #[mz_ore::instrument(level = "trace")]
1591 pub(super) fn drop_optimizer_notices(
1592 &mut self,
1593 dropped_entries: Vec<CatalogEntry>,
1594 ) -> BTreeSet<Arc<OptimizerNotice>> {
1595 let mut dropped_notices = BTreeSet::new();
1596 let mut drop_ids = BTreeSet::new();
1597
1598 for mut entry in dropped_entries {
1600 drop_ids.extend(entry.global_ids());
1601 if let Some(metainfo) = entry.item_mut().dataflow_metainfo_mut() {
1602 soft_assert_or_log!(
1603 metainfo.optimizer_notices.iter().all_unique(),
1604 "should have been pushed there by \
1605 `push_optimizer_notice_dedup`"
1606 );
1607 for n in metainfo.optimizer_notices.drain(..) {
1608 for dep_id in n.dependencies.iter() {
1611 if let Some(notices) = self.notices_by_dep_id.get_mut(dep_id) {
1612 notices.retain(|x| &n != x);
1613 if notices.is_empty() {
1614 self.notices_by_dep_id.remove(dep_id);
1615 }
1616 }
1617 }
1618 dropped_notices.insert(n);
1619 }
1620 }
1621 }
1622
1623 for id in &drop_ids {
1627 if let Some(notices) = self.notices_by_dep_id.remove(id) {
1628 for n in notices.into_iter() {
1629 if let Some(item_id) = n.item_id.as_ref() {
1634 if let Some(entry) = self.try_get_entry_by_global_id(item_id) {
1635 let catalog_item_id = entry.id();
1636 let entry = self.get_entry_mut(&catalog_item_id);
1637 if let Some(m) = entry.item_mut().dataflow_metainfo_mut() {
1638 m.optimizer_notices.retain(|x| &n != x);
1639 }
1640 }
1641 }
1642 dropped_notices.insert(n);
1643 }
1644 }
1645 }
1646
1647 let todo_dep_ids: BTreeSet<GlobalId> = dropped_notices
1650 .iter()
1651 .flat_map(|n| n.dependencies.iter())
1652 .filter(|dep_id| !drop_ids.contains(dep_id))
1653 .copied()
1654 .collect();
1655 for id in todo_dep_ids {
1656 if let Some(notices) = self.notices_by_dep_id.get_mut(&id) {
1657 notices.retain(|n| !dropped_notices.contains(n));
1658 if notices.is_empty() {
1659 self.notices_by_dep_id.remove(&id);
1660 }
1661 }
1662 }
1663
1664 dropped_notices
1665 }
1666
1667 fn get_schema_mut(
1668 &mut self,
1669 database_spec: &ResolvedDatabaseSpecifier,
1670 schema_spec: &SchemaSpecifier,
1671 conn_id: &ConnectionId,
1672 ) -> &mut Schema {
1673 match (database_spec, schema_spec) {
1675 (ResolvedDatabaseSpecifier::Ambient, SchemaSpecifier::Temporary) => self
1676 .temporary_namespaces
1677 .schema_mut(conn_id)
1678 .expect("catalog out of sync"),
1679 (ResolvedDatabaseSpecifier::Ambient, SchemaSpecifier::Id(id)) => self
1680 .ambient_schemas_by_id
1681 .get_mut(id)
1682 .expect("catalog out of sync"),
1683 (ResolvedDatabaseSpecifier::Id(database_id), SchemaSpecifier::Id(schema_id)) => self
1684 .database_by_id
1685 .get_mut(database_id)
1686 .expect("catalog out of sync")
1687 .schemas_by_id
1688 .get_mut(schema_id)
1689 .expect("catalog out of sync"),
1690 (ResolvedDatabaseSpecifier::Id(_), SchemaSpecifier::Temporary) => {
1691 unreachable!("temporary schemas are in the ambient database")
1692 }
1693 }
1694 }
1695
1696 #[instrument(name = "catalog::parse_views")]
1706 async fn parse_builtin_views(
1707 state: &mut CatalogState,
1708 builtin_views: Vec<(&'static BuiltinView, CatalogItemId, GlobalId)>,
1709 retractions: &mut InProgressRetractions,
1710 local_expression_cache: &mut LocalExpressionCache,
1711 ) -> Vec<BuiltinTableUpdate<&'static BuiltinTable>> {
1712 let mut builtin_table_updates = Vec::with_capacity(builtin_views.len());
1713 let (updates, additions): (Vec<_>, Vec<_>) =
1714 builtin_views
1715 .into_iter()
1716 .partition_map(|(view, item_id, gid)| {
1717 match retractions.system_object_mappings.remove(&item_id) {
1718 Some(entry) => Either::Left(entry),
1719 None => Either::Right((view, item_id, gid)),
1720 }
1721 });
1722
1723 for entry in updates {
1724 let item_id = entry.id();
1729 state.insert_entry(entry);
1730 builtin_table_updates.extend(state.pack_item_update(item_id, Diff::ONE));
1731 }
1732
1733 let mut handles = Vec::new();
1734 let mut awaiting_id_dependencies: BTreeMap<CatalogItemId, Vec<CatalogItemId>> =
1735 BTreeMap::new();
1736 let mut awaiting_name_dependencies: BTreeMap<String, Vec<CatalogItemId>> = BTreeMap::new();
1737 let mut awaiting_all = Vec::new();
1740 let mut completed_ids: BTreeSet<CatalogItemId> = BTreeSet::new();
1742 let mut completed_names: BTreeSet<String> = BTreeSet::new();
1743
1744 let mut views: BTreeMap<CatalogItemId, (&BuiltinView, GlobalId)> = additions
1746 .into_iter()
1747 .map(|(view, item_id, gid)| (item_id, (view, gid)))
1748 .collect();
1749 let item_ids: Vec<_> = views.keys().copied().collect();
1750
1751 let mut ready: VecDeque<CatalogItemId> = views.keys().cloned().collect();
1752 while !handles.is_empty() || !ready.is_empty() || !awaiting_all.is_empty() {
1753 if handles.is_empty() && ready.is_empty() {
1754 ready.extend(awaiting_all.drain(..));
1756 }
1757
1758 if !ready.is_empty() {
1760 let spawn_state = Arc::new(state.clone());
1761 while let Some(id) = ready.pop_front() {
1762 let (view, global_id) = views.get(&id).expect("must exist");
1763 let global_id = *global_id;
1764 let create_sql = view.create_sql();
1765 let versions = BTreeMap::new();
1767
1768 let span = info_span!(parent: None, "parse builtin view", name = view.name);
1769 OpenTelemetryContext::obtain().attach_as_parent_to(&span);
1770 let task_state = Arc::clone(&spawn_state);
1771 let cached_expr = local_expression_cache.remove_cached_expression(&global_id);
1772 let handle = mz_ore::task::spawn_blocking(
1773 || "parse view",
1774 move || {
1775 span.in_scope(|| {
1776 let res = task_state.parse_item_inner(
1777 global_id,
1778 &create_sql,
1779 &versions,
1780 None,
1781 false,
1782 None,
1783 cached_expr,
1784 None,
1785 );
1786 (id, global_id, res)
1787 })
1788 },
1789 );
1790 handles.push(handle);
1791 }
1792 }
1793
1794 let (selected, _idx, remaining) = future::select_all(handles).await;
1796 handles = remaining;
1797 let (id, global_id, res) = selected;
1798 let mut insert_cached_expr = |cached_expr| {
1799 if let Some(cached_expr) = cached_expr {
1800 local_expression_cache.insert_cached_expression(global_id, cached_expr);
1801 }
1802 };
1803 match res {
1804 Ok((item, uncached_expr)) => {
1805 if let Some((uncached_expr, optimizer_features)) = uncached_expr {
1806 local_expression_cache.insert_uncached_expression(
1807 global_id,
1808 uncached_expr,
1809 optimizer_features,
1810 RelationVersion::root(),
1811 );
1812 }
1813 let (view, _gid) = views.remove(&id).expect("must exist");
1815 let schema_id = state
1816 .ambient_schemas_by_name
1817 .get(view.schema)
1818 .unwrap_or_else(|| panic!("unknown ambient schema: {}", view.schema));
1819 let qname = QualifiedItemName {
1820 qualifiers: ItemQualifiers {
1821 database_spec: ResolvedDatabaseSpecifier::Ambient,
1822 schema_spec: SchemaSpecifier::Id(*schema_id),
1823 },
1824 item: view.name.into(),
1825 };
1826 let mut acl_items = vec![rbac::owner_privilege(
1827 mz_sql::catalog::ObjectType::View,
1828 MZ_SYSTEM_ROLE_ID,
1829 )];
1830 acl_items.extend_from_slice(&view.access);
1831
1832 state.insert_item(
1833 id,
1834 view.oid,
1835 qname,
1836 item,
1837 MZ_SYSTEM_ROLE_ID,
1838 PrivilegeMap::from_mz_acl_items(acl_items),
1839 );
1840
1841 let mut resolved_dependent_items = Vec::new();
1843 if let Some(dependent_items) = awaiting_id_dependencies.remove(&id) {
1844 resolved_dependent_items.extend(dependent_items);
1845 }
1846 let entry = state.get_entry(&id);
1847 let full_name = state.resolve_full_name(entry.name(), None).to_string();
1848 if let Some(dependent_items) = awaiting_name_dependencies.remove(&full_name) {
1849 resolved_dependent_items.extend(dependent_items);
1850 }
1851 ready.extend(resolved_dependent_items);
1852
1853 completed_ids.insert(id);
1854 completed_names.insert(full_name);
1855 }
1856 Err((
1858 AdapterError::PlanError(plan::PlanError::InvalidId(missing_dep)),
1859 cached_expr,
1860 )) => {
1861 insert_cached_expr(cached_expr);
1862 if completed_ids.contains(&missing_dep) {
1863 ready.push_back(id);
1864 } else {
1865 awaiting_id_dependencies
1866 .entry(missing_dep)
1867 .or_default()
1868 .push(id);
1869 }
1870 }
1871 Err((
1873 AdapterError::PlanError(plan::PlanError::Catalog(
1874 SqlCatalogError::UnknownItem(missing_dep),
1875 )),
1876 cached_expr,
1877 )) => {
1878 insert_cached_expr(cached_expr);
1879 match CatalogItemId::from_str(&missing_dep) {
1880 Ok(missing_dep) => {
1881 if completed_ids.contains(&missing_dep) {
1882 ready.push_back(id);
1883 } else {
1884 awaiting_id_dependencies
1885 .entry(missing_dep)
1886 .or_default()
1887 .push(id);
1888 }
1889 }
1890 Err(_) => {
1891 if completed_names.contains(&missing_dep) {
1892 ready.push_back(id);
1893 } else {
1894 awaiting_name_dependencies
1895 .entry(missing_dep)
1896 .or_default()
1897 .push(id);
1898 }
1899 }
1900 }
1901 }
1902 Err((
1903 AdapterError::PlanError(plan::PlanError::InvalidCast { .. }),
1904 cached_expr,
1905 )) => {
1906 insert_cached_expr(cached_expr);
1907 awaiting_all.push(id);
1908 }
1909 Err((e, _)) => {
1910 let (bad_view, _gid) = views.get(&id).expect("must exist");
1911 panic!(
1912 "internal error: failed to load bootstrap view:\n\
1913 {name}\n\
1914 error:\n\
1915 {e:?}\n\n\
1916 Make sure that the schema name is specified in the builtin view's create sql statement.
1917 ",
1918 name = bad_view.name,
1919 )
1920 }
1921 }
1922 }
1923
1924 assert!(awaiting_id_dependencies.is_empty());
1925 assert!(
1926 awaiting_name_dependencies.is_empty(),
1927 "awaiting_name_dependencies: {awaiting_name_dependencies:?}"
1928 );
1929 assert!(awaiting_all.is_empty());
1930 assert!(views.is_empty());
1931
1932 builtin_table_updates.extend(
1934 item_ids
1935 .into_iter()
1936 .flat_map(|id| state.pack_item_update(id, Diff::ONE)),
1937 );
1938
1939 builtin_table_updates
1940 }
1941
1942 fn insert_entry(&mut self, entry: CatalogEntry) {
1944 if !entry.id.is_system() {
1945 if let Some(cluster_id) = entry.item.cluster_id() {
1946 self.clusters_by_id
1947 .get_mut(&cluster_id)
1948 .expect("catalog out of sync")
1949 .bound_objects
1950 .insert(entry.id);
1951 };
1952 }
1953
1954 for u in entry.references().items() {
1955 match self.entry_by_id.get_mut(u) {
1956 Some(metadata) => metadata.referenced_by.push(entry.id()),
1957 None => panic!(
1958 "Catalog: missing dependent catalog item {} while installing {}",
1959 u,
1960 self.resolve_full_name(entry.name(), entry.conn_id())
1961 ),
1962 }
1963 }
1964 for u in entry.uses() {
1965 if u == entry.id() {
1968 continue;
1969 }
1970 match self.entry_by_id.get_mut(&u) {
1971 Some(metadata) => metadata.used_by.push(entry.id()),
1972 None => panic!(
1973 "Catalog: missing dependent catalog item {} while installing {}",
1974 u,
1975 self.resolve_full_name(entry.name(), entry.conn_id())
1976 ),
1977 }
1978 }
1979 for gid in entry.item.global_ids() {
1980 self.entry_by_global_id.insert(gid, entry.id());
1981 }
1982 let conn_id = entry.item().conn_id().unwrap_or(&SYSTEM_CONN_ID);
1983 if entry.name().qualifiers.schema_spec == SchemaSpecifier::Temporary {
1986 self.temporary_namespaces
1987 .ensure_schema(conn_id, entry.owner_id);
1988 }
1989 let schema = self.get_schema_mut(
1990 &entry.name().qualifiers.database_spec,
1991 &entry.name().qualifiers.schema_spec,
1992 conn_id,
1993 );
1994
1995 let prev_id = match entry.item() {
1996 CatalogItem::Func(_) => schema
1997 .functions
1998 .insert(entry.name().item.clone(), entry.id()),
1999 CatalogItem::Type(_) => schema.types.insert(entry.name().item.clone(), entry.id()),
2000 _ => schema.items.insert(entry.name().item.clone(), entry.id()),
2001 };
2002
2003 assert!(
2004 prev_id.is_none(),
2005 "builtin name collision on {:?}",
2006 entry.name().item.clone()
2007 );
2008
2009 self.entry_by_id.insert(entry.id(), entry.clone());
2010 }
2011
2012 fn insert_item(
2014 &mut self,
2015 id: CatalogItemId,
2016 oid: u32,
2017 name: QualifiedItemName,
2018 item: CatalogItem,
2019 owner_id: RoleId,
2020 privileges: PrivilegeMap,
2021 ) {
2022 let entry = CatalogEntry {
2023 item,
2024 name,
2025 id,
2026 oid,
2027 used_by: Vec::new(),
2028 referenced_by: Vec::new(),
2029 owner_id,
2030 privileges,
2031 };
2032
2033 self.insert_entry(entry);
2034 }
2035
2036 #[mz_ore::instrument(level = "trace")]
2037 fn drop_item(&mut self, id: CatalogItemId) -> CatalogEntry {
2038 let metadata = self.entry_by_id.remove(&id).expect("catalog out of sync");
2039 for u in metadata.references().items() {
2040 if let Some(dep_metadata) = self.entry_by_id.get_mut(u) {
2041 dep_metadata.referenced_by.retain(|u| *u != metadata.id())
2042 }
2043 }
2044 for u in metadata.uses() {
2045 if let Some(dep_metadata) = self.entry_by_id.get_mut(&u) {
2046 dep_metadata.used_by.retain(|u| *u != metadata.id())
2047 }
2048 }
2049 for gid in metadata.global_ids() {
2050 self.entry_by_global_id.remove(&gid);
2051 }
2052
2053 let conn_id = metadata.item().conn_id().unwrap_or(&SYSTEM_CONN_ID);
2054 let schema = self.get_schema_mut(
2055 &metadata.name().qualifiers.database_spec,
2056 &metadata.name().qualifiers.schema_spec,
2057 conn_id,
2058 );
2059 if metadata.item_type() == CatalogItemType::Type {
2060 schema
2061 .types
2062 .remove(&metadata.name().item)
2063 .expect("catalog out of sync");
2064 } else {
2065 assert_ne!(metadata.item_type(), CatalogItemType::Func);
2068
2069 schema
2070 .items
2071 .remove(&metadata.name().item)
2072 .expect("catalog out of sync");
2073 };
2074
2075 if !id.is_system() {
2076 if let Some(cluster_id) = metadata.item().cluster_id() {
2077 assert!(
2078 self.clusters_by_id
2079 .get_mut(&cluster_id)
2080 .expect("catalog out of sync")
2081 .bound_objects
2082 .remove(&id),
2083 "catalog out of sync"
2084 );
2085 }
2086 }
2087
2088 metadata
2089 }
2090
2091 fn insert_introspection_source_index(
2092 &mut self,
2093 cluster_id: ClusterId,
2094 log: &'static BuiltinLog,
2095 item_id: CatalogItemId,
2096 global_id: GlobalId,
2097 oid: u32,
2098 ) {
2099 let (index_name, index) =
2100 self.create_introspection_source_index(cluster_id, log, global_id);
2101 self.insert_item(
2102 item_id,
2103 oid,
2104 index_name,
2105 index,
2106 MZ_SYSTEM_ROLE_ID,
2107 PrivilegeMap::default(),
2108 );
2109 }
2110
2111 fn create_introspection_source_index(
2112 &self,
2113 cluster_id: ClusterId,
2114 log: &'static BuiltinLog,
2115 global_id: GlobalId,
2116 ) -> (QualifiedItemName, CatalogItem) {
2117 let source_name = FullItemName {
2118 database: RawDatabaseSpecifier::Ambient,
2119 schema: log.schema.into(),
2120 item: log.name.into(),
2121 };
2122 let index_name = format!("{}_{}_primary_idx", log.name, cluster_id);
2123 let mut index_name = QualifiedItemName {
2124 qualifiers: ItemQualifiers {
2125 database_spec: ResolvedDatabaseSpecifier::Ambient,
2126 schema_spec: SchemaSpecifier::Id(self.get_mz_introspection_schema_id()),
2127 },
2128 item: index_name.clone(),
2129 };
2130 index_name = self.find_available_name(index_name, &SYSTEM_CONN_ID);
2131 let index_item_name = index_name.item.clone();
2132 let (log_item_id, log_global_id) = self.resolve_builtin_log(log);
2133 let index = CatalogItem::Index(Index {
2134 global_id,
2135 on: log_global_id,
2136 keys: log
2137 .variant
2138 .index_by()
2139 .into_iter()
2140 .map(MirScalarExpr::column)
2141 .collect(),
2142 create_sql: index_sql(
2143 index_item_name,
2144 cluster_id,
2145 source_name,
2146 &log.variant.desc(),
2147 &log.variant.index_by(),
2148 ),
2149 conn_id: None,
2150 resolved_ids: [(log_item_id, log_global_id)].into_iter().collect(),
2151 cluster_id,
2152 is_retained_metrics_object: false,
2153 custom_logical_compaction_window: None,
2154 optimized_plan: None,
2155 physical_plan: None,
2156 dataflow_metainfo: None,
2157 });
2158 (index_name, index)
2159 }
2160
2161 fn insert_system_configuration(&mut self, name: &str, value: VarInput) -> Result<bool, Error> {
2166 Ok(Arc::make_mut(&mut self.system_configuration).set(name, value)?)
2167 }
2168
2169 fn remove_system_configuration(&mut self, name: &str) -> Result<bool, Error> {
2174 Ok(Arc::make_mut(&mut self.system_configuration).reset(name)?)
2175 }
2176}
2177
2178fn sort_updates(updates: Vec<StateUpdate>) -> Vec<StateUpdate> {
2186 fn push_update<T>(
2187 update: T,
2188 diff: StateDiff,
2189 retractions: &mut Vec<T>,
2190 additions: &mut Vec<T>,
2191 ) {
2192 match diff {
2193 StateDiff::Retraction => retractions.push(update),
2194 StateDiff::Addition => additions.push(update),
2195 }
2196 }
2197
2198 soft_assert_no_log!(
2199 updates.iter().map(|update| update.ts).all_equal(),
2200 "all timestamps should be equal: {updates:?}"
2201 );
2202 soft_assert_no_log!(
2203 {
2204 let mut dedup = BTreeSet::new();
2205 updates.iter().all(|update| dedup.insert(&update.kind))
2206 },
2207 "updates should be consolidated: {updates:?}"
2208 );
2209
2210 let mut pre_cluster_retractions = Vec::new();
2212 let mut pre_cluster_additions = Vec::new();
2213 let mut cluster_retractions = Vec::new();
2214 let mut cluster_additions = Vec::new();
2215 let mut builtin_item_updates = Vec::new();
2216 let mut item_retractions = Vec::new();
2217 let mut item_additions = Vec::new();
2218 let mut post_item_retractions = Vec::new();
2219 let mut post_item_additions = Vec::new();
2220 for update in updates {
2221 let diff = update.diff.clone();
2222 match update.kind {
2223 StateUpdateKind::Role(_)
2224 | StateUpdateKind::RoleAuth(_)
2225 | StateUpdateKind::Database(_)
2226 | StateUpdateKind::Schema(_)
2227 | StateUpdateKind::DefaultPrivilege(_)
2228 | StateUpdateKind::SystemPrivilege(_)
2229 | StateUpdateKind::SystemConfiguration(_)
2230 | StateUpdateKind::NetworkPolicy(_) => push_update(
2231 update,
2232 diff,
2233 &mut pre_cluster_retractions,
2234 &mut pre_cluster_additions,
2235 ),
2236 StateUpdateKind::Cluster(_)
2237 | StateUpdateKind::ClusterSystemConfiguration(_)
2238 | StateUpdateKind::IntrospectionSourceIndex(_)
2239 | StateUpdateKind::ClusterReplica(_)
2240 | StateUpdateKind::ReplicaSystemConfiguration(_) => push_update(
2241 update,
2242 diff,
2243 &mut cluster_retractions,
2244 &mut cluster_additions,
2245 ),
2246 StateUpdateKind::SystemObjectMapping(system_object_mapping) => {
2247 builtin_item_updates.push((system_object_mapping, update.ts, update.diff))
2248 }
2249 StateUpdateKind::Item(item) => push_update(
2250 (item, update.ts, update.diff),
2251 diff,
2252 &mut item_retractions,
2253 &mut item_additions,
2254 ),
2255 StateUpdateKind::Comment(_)
2256 | StateUpdateKind::SourceReferences(_)
2257 | StateUpdateKind::AuditLog(_)
2258 | StateUpdateKind::StorageCollectionMetadata(_)
2259 | StateUpdateKind::UnfinalizedShard(_) => push_update(
2260 update,
2261 diff,
2262 &mut post_item_retractions,
2263 &mut post_item_additions,
2264 ),
2265 }
2266 }
2267
2268 let builtin_item_updates = builtin_item_updates
2271 .into_iter()
2272 .map(|(system_object_mapping, ts, diff)| {
2273 let idx = BUILTIN_LOOKUP
2274 .get(&system_object_mapping.description)
2275 .expect("missing builtin")
2276 .0;
2277 (idx, system_object_mapping, ts, diff)
2278 })
2279 .sorted_by_key(|(idx, _, _, _)| *idx)
2280 .map(|(_, system_object_mapping, ts, diff)| (system_object_mapping, ts, diff));
2281
2282 let mut builtin_source_retractions = Vec::new();
2286 let mut builtin_source_additions = Vec::new();
2287 let mut other_builtin_retractions = Vec::new();
2288 let mut other_builtin_additions = Vec::new();
2289 for (builtin_item_update, ts, diff) in builtin_item_updates {
2290 let object_type = builtin_item_update.description.object_type;
2291 let update = StateUpdate {
2292 kind: StateUpdateKind::SystemObjectMapping(builtin_item_update),
2293 ts,
2294 diff,
2295 };
2296 if object_type == CatalogItemType::Source {
2297 push_update(
2298 update,
2299 diff,
2300 &mut builtin_source_retractions,
2301 &mut builtin_source_additions,
2302 );
2303 } else {
2304 push_update(
2305 update,
2306 diff,
2307 &mut other_builtin_retractions,
2308 &mut other_builtin_additions,
2309 );
2310 }
2311 }
2312
2313 fn sort_items_topological(items: &mut Vec<(mz_catalog::durable::Item, Timestamp, StateDiff)>) {
2319 tracing::debug!(?items, "sorting items by dependencies");
2320
2321 let key_fn = |item: &(mz_catalog::durable::Item, _, _)| item.0.id;
2322 let dependencies_fn = |item: &(mz_catalog::durable::Item, _, _)| {
2323 let statement = mz_sql::parse::parse(&item.0.create_sql)
2324 .expect("valid create_sql")
2325 .into_element()
2326 .ast;
2327 mz_sql::names::dependencies(&statement).expect("failed to find dependencies of item")
2328 };
2329 sort_topological(items, key_fn, dependencies_fn);
2330 }
2331
2332 fn sort_item_updates(
2348 item_updates: Vec<(mz_catalog::durable::Item, Timestamp, StateDiff)>,
2349 ) -> VecDeque<(mz_catalog::durable::Item, Timestamp, StateDiff)> {
2350 let mut types = Vec::new();
2353 let mut funcs = Vec::new();
2356 let mut secrets = Vec::new();
2357 let mut connections = Vec::new();
2358 let mut sources = Vec::new();
2359 let mut tables = Vec::new();
2360 let mut derived_items = Vec::new();
2361 let mut sinks = Vec::new();
2362 for update in item_updates {
2363 match update.0.item_type() {
2364 CatalogItemType::Type => types.push(update),
2365 CatalogItemType::Func => funcs.push(update),
2366 CatalogItemType::Secret => secrets.push(update),
2367 CatalogItemType::Connection => connections.push(update),
2368 CatalogItemType::Source => sources.push(update),
2369 CatalogItemType::Table => tables.push(update),
2370 CatalogItemType::View
2371 | CatalogItemType::MaterializedView
2372 | CatalogItemType::Index
2373 | CatalogItemType::MetricSink => derived_items.push(update),
2374 CatalogItemType::Sink => sinks.push(update),
2375 }
2376 }
2377
2378 sort_items_topological(&mut connections);
2382 sort_items_topological(&mut derived_items);
2383
2384 for group in [
2386 &mut types,
2387 &mut funcs,
2388 &mut secrets,
2389 &mut sources,
2390 &mut tables,
2391 &mut sinks,
2392 ] {
2393 group.sort_by_key(|(item, _, _)| item.id);
2394 }
2395
2396 iter::empty()
2397 .chain(types)
2398 .chain(funcs)
2399 .chain(secrets)
2400 .chain(connections)
2401 .chain(sources)
2402 .chain(tables)
2403 .chain(derived_items)
2404 .chain(sinks)
2405 .collect()
2406 }
2407
2408 fn into_state_updates(
2412 item_updates: VecDeque<(mz_catalog::durable::Item, Timestamp, StateDiff)>,
2413 ) -> Vec<StateUpdate> {
2414 item_updates
2415 .into_iter()
2416 .map(|(item, ts, diff)| StateUpdate {
2417 kind: StateUpdateKind::Item(item),
2418 ts,
2419 diff,
2420 })
2421 .collect()
2422 }
2423 let item_retractions = into_state_updates(sort_item_updates(item_retractions));
2424 let item_additions = into_state_updates(sort_item_updates(item_additions));
2425
2426 iter::empty()
2428 .chain(post_item_retractions.into_iter().rev())
2430 .chain(item_retractions.into_iter().rev())
2431 .chain(other_builtin_retractions.into_iter().rev())
2432 .chain(cluster_retractions.into_iter().rev())
2433 .chain(builtin_source_retractions.into_iter().rev())
2434 .chain(pre_cluster_retractions.into_iter().rev())
2435 .chain(pre_cluster_additions)
2436 .chain(builtin_source_additions)
2437 .chain(cluster_additions)
2438 .chain(other_builtin_additions)
2439 .chain(item_additions)
2440 .chain(post_item_additions)
2441 .collect()
2442}
2443
2444enum ApplyState {
2449 BuiltinViewAdditions(Vec<(&'static BuiltinView, CatalogItemId, GlobalId)>),
2451 Items(Vec<StateUpdate>),
2457 Updates(Vec<StateUpdate>),
2459}
2460
2461impl ApplyState {
2462 fn new(update: StateUpdate) -> Self {
2463 use StateUpdateKind::*;
2464 match &update.kind {
2465 SystemObjectMapping(som)
2466 if som.description.object_type == CatalogItemType::View
2467 && update.diff == StateDiff::Addition =>
2468 {
2469 let view_addition = lookup_builtin_view_addition(som.clone());
2470 Self::BuiltinViewAdditions(vec![view_addition])
2471 }
2472
2473 IntrospectionSourceIndex(_) | SystemObjectMapping(_) | Item(_) => {
2474 Self::Items(vec![update])
2475 }
2476
2477 Role(_)
2478 | RoleAuth(_)
2479 | Database(_)
2480 | Schema(_)
2481 | DefaultPrivilege(_)
2482 | SystemPrivilege(_)
2483 | SystemConfiguration(_)
2484 | ClusterSystemConfiguration(_)
2485 | ReplicaSystemConfiguration(_)
2486 | Cluster(_)
2487 | NetworkPolicy(_)
2488 | ClusterReplica(_)
2489 | SourceReferences(_)
2490 | Comment(_)
2491 | AuditLog(_)
2492 | StorageCollectionMetadata(_)
2493 | UnfinalizedShard(_) => Self::Updates(vec![update]),
2494 }
2495 }
2496
2497 async fn apply(
2503 self,
2504 state: &mut CatalogState,
2505 retractions: &mut InProgressRetractions,
2506 local_expression_cache: &mut LocalExpressionCache,
2507 ) -> (
2508 Vec<BuiltinTableUpdate<&'static BuiltinTable>>,
2509 Vec<ParsedStateUpdate>,
2510 ) {
2511 match self {
2512 Self::BuiltinViewAdditions(builtin_view_additions) => {
2513 let restore = Arc::clone(&state.system_configuration);
2514 Arc::make_mut(&mut state.system_configuration).enable_for_item_parsing();
2515 let builtin_table_updates = CatalogState::parse_builtin_views(
2516 state,
2517 builtin_view_additions,
2518 retractions,
2519 local_expression_cache,
2520 )
2521 .await;
2522 state.system_configuration = restore;
2523 (builtin_table_updates, Vec::new())
2524 }
2525 Self::Items(updates) => state.with_enable_for_item_parsing(|state| {
2526 state
2527 .apply_updates_inner(updates, retractions, local_expression_cache)
2528 .expect("corrupt catalog")
2529 }),
2530 Self::Updates(updates) => state
2531 .apply_updates_inner(updates, retractions, local_expression_cache)
2532 .expect("corrupt catalog"),
2533 }
2534 }
2535
2536 async fn step(
2537 self,
2538 next: Self,
2539 state: &mut CatalogState,
2540 retractions: &mut InProgressRetractions,
2541 local_expression_cache: &mut LocalExpressionCache,
2542 ) -> (
2543 Self,
2544 (
2545 Vec<BuiltinTableUpdate<&'static BuiltinTable>>,
2546 Vec<ParsedStateUpdate>,
2547 ),
2548 ) {
2549 match (self, next) {
2550 (
2551 Self::BuiltinViewAdditions(mut builtin_view_additions),
2552 Self::BuiltinViewAdditions(next_builtin_view_additions),
2553 ) => {
2554 builtin_view_additions.extend(next_builtin_view_additions);
2556 (
2557 Self::BuiltinViewAdditions(builtin_view_additions),
2558 (Vec::new(), Vec::new()),
2559 )
2560 }
2561 (Self::Items(mut updates), Self::Items(next_updates)) => {
2562 updates.extend(next_updates);
2564 (Self::Items(updates), (Vec::new(), Vec::new()))
2565 }
2566 (Self::Updates(mut updates), Self::Updates(next_updates)) => {
2567 updates.extend(next_updates);
2569 (Self::Updates(updates), (Vec::new(), Vec::new()))
2570 }
2571 (apply_state, next_apply_state) => {
2572 let updates = apply_state
2574 .apply(state, retractions, local_expression_cache)
2575 .await;
2576 (next_apply_state, updates)
2577 }
2578 }
2579 }
2580}
2581
2582trait MutableMap<K, V> {
2585 fn insert(&mut self, key: K, value: V) -> Option<V>;
2586 fn remove(&mut self, key: &K) -> Option<V>;
2587}
2588
2589impl<K: Ord, V> MutableMap<K, V> for BTreeMap<K, V> {
2590 fn insert(&mut self, key: K, value: V) -> Option<V> {
2591 BTreeMap::insert(self, key, value)
2592 }
2593 fn remove(&mut self, key: &K) -> Option<V> {
2594 BTreeMap::remove(self, key)
2595 }
2596}
2597
2598impl<K: Ord + Clone, V: Clone> MutableMap<K, V> for imbl::OrdMap<K, V> {
2599 fn insert(&mut self, key: K, value: V) -> Option<V> {
2600 imbl::OrdMap::insert(self, key, value)
2601 }
2602 fn remove(&mut self, key: &K) -> Option<V> {
2603 imbl::OrdMap::remove(self, key)
2604 }
2605}
2606
2607fn apply_inverted_lookup<K, V>(map: &mut impl MutableMap<K, V>, key: &K, value: V, diff: StateDiff)
2612where
2613 K: Ord + Clone + Debug,
2614 V: PartialEq + Debug,
2615{
2616 match diff {
2617 StateDiff::Retraction => {
2618 let prev = map.remove(key);
2619 assert_eq!(
2620 prev,
2621 Some(value),
2622 "retraction does not match existing value: {key:?}"
2623 );
2624 }
2625 StateDiff::Addition => {
2626 let prev = map.insert(key.clone(), value);
2627 assert_eq!(
2628 prev, None,
2629 "values must be explicitly retracted before inserting a new value: {key:?}"
2630 );
2631 }
2632 }
2633}
2634
2635fn apply_with_update<K, V, D>(
2638 map: &mut impl MutableMap<K, V>,
2639 durable: D,
2640 key_fn: impl FnOnce(&D) -> K,
2641 diff: StateDiff,
2642 retractions: &mut BTreeMap<D::Key, V>,
2643) where
2644 K: Ord,
2645 V: UpdateFrom<D> + PartialEq + Debug,
2646 D: DurableType,
2647 D::Key: Ord,
2648{
2649 match diff {
2650 StateDiff::Retraction => {
2651 let mem_key = key_fn(&durable);
2652 let value = map
2653 .remove(&mem_key)
2654 .expect("retraction does not match existing value: {key:?}");
2655 let durable_key = durable.into_key_value().0;
2656 retractions.insert(durable_key, value);
2657 }
2658 StateDiff::Addition => {
2659 let mem_key = key_fn(&durable);
2660 let durable_key = durable.key();
2661 let value = match retractions.remove(&durable_key) {
2662 Some(mut retraction) => {
2663 retraction.update_from(durable);
2664 retraction
2665 }
2666 None => durable.into(),
2667 };
2668 let prev = map.insert(mem_key, value);
2669 assert_eq!(
2670 prev, None,
2671 "values must be explicitly retracted before inserting a new value"
2672 );
2673 }
2674 }
2675}
2676
2677fn lookup_builtin_view_addition(
2679 mapping: SystemObjectMapping,
2680) -> (&'static BuiltinView, CatalogItemId, GlobalId) {
2681 let (_, builtin) = BUILTIN_LOOKUP
2682 .get(&mapping.description)
2683 .expect("missing builtin view");
2684 let Builtin::View(view) = builtin else {
2685 unreachable!("programming error, expected BuiltinView found {builtin:?}");
2686 };
2687
2688 (
2689 view,
2690 mapping.unique_identifier.catalog_id,
2691 mapping.unique_identifier.global_id,
2692 )
2693}