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(_) => Vec::new(),
1486 StateUpdateKind::AuditLog(_) => Vec::new(),
1490 StateUpdateKind::Database(_)
1491 | StateUpdateKind::Schema(_)
1492 | StateUpdateKind::NetworkPolicy(_)
1493 | StateUpdateKind::StorageCollectionMetadata(_)
1494 | StateUpdateKind::UnfinalizedShard(_) => Vec::new(),
1495 }
1496 }
1497
1498 fn get_entry_mut(&mut self, id: &CatalogItemId) -> &mut CatalogEntry {
1499 self.entry_by_id
1500 .get_mut(id)
1501 .unwrap_or_else(|| panic!("catalog out of sync, missing id {id}"))
1502 }
1503
1504 pub(super) fn set_optimized_plan(
1509 &mut self,
1510 id: GlobalId,
1511 plan: DataflowDescription<mz_expr::OptimizedMirRelationExpr>,
1512 ) {
1513 let item_id = self.entry_by_global_id[&id];
1514 let entry = self.get_entry_mut(&item_id);
1515 match entry.item_mut() {
1516 CatalogItem::Index(idx) => idx.optimized_plan = Some(Arc::new(plan)),
1517 CatalogItem::MaterializedView(mv) => mv.optimized_plan = Some(Arc::new(plan)),
1518 CatalogItem::MetricSink(ms) => ms.optimized_plan = Some(Arc::new(plan)),
1519 other => panic!("set_optimized_plan called on {} ({:?})", id, other.typ()),
1520 }
1521 }
1522
1523 pub(super) fn set_physical_plan(
1528 &mut self,
1529 id: GlobalId,
1530 plan: DataflowDescription<mz_compute_types::plan::LirRelationExpr>,
1531 ) {
1532 let item_id = self.entry_by_global_id[&id];
1533 let entry = self.get_entry_mut(&item_id);
1534 match entry.item_mut() {
1535 CatalogItem::Index(idx) => idx.physical_plan = Some(Arc::new(plan)),
1536 CatalogItem::MaterializedView(mv) => mv.physical_plan = Some(Arc::new(plan)),
1537 CatalogItem::MetricSink(ms) => ms.physical_plan = Some(Arc::new(plan)),
1538 other => panic!("set_physical_plan called on {} ({:?})", id, other.typ()),
1539 }
1540 }
1541
1542 pub(super) fn set_dataflow_metainfo(
1547 &mut self,
1548 id: GlobalId,
1549 metainfo: DataflowMetainfo<Arc<OptimizerNotice>>,
1550 ) {
1551 for notice in metainfo.optimizer_notices.iter() {
1553 for dep_id in notice.dependencies.iter() {
1554 self.notices_by_dep_id
1555 .entry(*dep_id)
1556 .or_default()
1557 .push(Arc::clone(notice));
1558 }
1559 if let Some(item_id) = notice.item_id {
1560 soft_assert_eq_or_log!(
1561 item_id,
1562 id,
1563 "notice.item_id should match the id for whom we are saving the notice"
1564 );
1565 }
1566 }
1567 let item_id = self.entry_by_global_id[&id];
1569 let entry = self.get_entry_mut(&item_id);
1570 match entry.item_mut() {
1571 CatalogItem::Index(idx) => idx.dataflow_metainfo = Some(metainfo),
1572 CatalogItem::MaterializedView(mv) => mv.dataflow_metainfo = Some(metainfo),
1573 CatalogItem::MetricSink(ms) => ms.dataflow_metainfo = Some(metainfo),
1574 other => panic!("set_dataflow_metainfo called on {} ({:?})", id, other.typ()),
1575 }
1576 }
1577
1578 #[mz_ore::instrument(level = "trace")]
1589 pub(super) fn drop_optimizer_notices(
1590 &mut self,
1591 dropped_entries: Vec<CatalogEntry>,
1592 ) -> BTreeSet<Arc<OptimizerNotice>> {
1593 let mut dropped_notices = BTreeSet::new();
1594 let mut drop_ids = BTreeSet::new();
1595
1596 for mut entry in dropped_entries {
1598 drop_ids.extend(entry.global_ids());
1599 if let Some(metainfo) = entry.item_mut().dataflow_metainfo_mut() {
1600 soft_assert_or_log!(
1601 metainfo.optimizer_notices.iter().all_unique(),
1602 "should have been pushed there by \
1603 `push_optimizer_notice_dedup`"
1604 );
1605 for n in metainfo.optimizer_notices.drain(..) {
1606 for dep_id in n.dependencies.iter() {
1609 if let Some(notices) = self.notices_by_dep_id.get_mut(dep_id) {
1610 notices.retain(|x| &n != x);
1611 if notices.is_empty() {
1612 self.notices_by_dep_id.remove(dep_id);
1613 }
1614 }
1615 }
1616 dropped_notices.insert(n);
1617 }
1618 }
1619 }
1620
1621 for id in &drop_ids {
1625 if let Some(notices) = self.notices_by_dep_id.remove(id) {
1626 for n in notices.into_iter() {
1627 if let Some(item_id) = n.item_id.as_ref() {
1632 if let Some(entry) = self.try_get_entry_by_global_id(item_id) {
1633 let catalog_item_id = entry.id();
1634 let entry = self.get_entry_mut(&catalog_item_id);
1635 if let Some(m) = entry.item_mut().dataflow_metainfo_mut() {
1636 m.optimizer_notices.retain(|x| &n != x);
1637 }
1638 }
1639 }
1640 dropped_notices.insert(n);
1641 }
1642 }
1643 }
1644
1645 let todo_dep_ids: BTreeSet<GlobalId> = dropped_notices
1648 .iter()
1649 .flat_map(|n| n.dependencies.iter())
1650 .filter(|dep_id| !drop_ids.contains(dep_id))
1651 .copied()
1652 .collect();
1653 for id in todo_dep_ids {
1654 if let Some(notices) = self.notices_by_dep_id.get_mut(&id) {
1655 notices.retain(|n| !dropped_notices.contains(n));
1656 if notices.is_empty() {
1657 self.notices_by_dep_id.remove(&id);
1658 }
1659 }
1660 }
1661
1662 dropped_notices
1663 }
1664
1665 fn get_schema_mut(
1666 &mut self,
1667 database_spec: &ResolvedDatabaseSpecifier,
1668 schema_spec: &SchemaSpecifier,
1669 conn_id: &ConnectionId,
1670 ) -> &mut Schema {
1671 match (database_spec, schema_spec) {
1673 (ResolvedDatabaseSpecifier::Ambient, SchemaSpecifier::Temporary) => self
1674 .temporary_namespaces
1675 .schema_mut(conn_id)
1676 .expect("catalog out of sync"),
1677 (ResolvedDatabaseSpecifier::Ambient, SchemaSpecifier::Id(id)) => self
1678 .ambient_schemas_by_id
1679 .get_mut(id)
1680 .expect("catalog out of sync"),
1681 (ResolvedDatabaseSpecifier::Id(database_id), SchemaSpecifier::Id(schema_id)) => self
1682 .database_by_id
1683 .get_mut(database_id)
1684 .expect("catalog out of sync")
1685 .schemas_by_id
1686 .get_mut(schema_id)
1687 .expect("catalog out of sync"),
1688 (ResolvedDatabaseSpecifier::Id(_), SchemaSpecifier::Temporary) => {
1689 unreachable!("temporary schemas are in the ambient database")
1690 }
1691 }
1692 }
1693
1694 #[instrument(name = "catalog::parse_views")]
1704 async fn parse_builtin_views(
1705 state: &mut CatalogState,
1706 builtin_views: Vec<(&'static BuiltinView, CatalogItemId, GlobalId)>,
1707 retractions: &mut InProgressRetractions,
1708 local_expression_cache: &mut LocalExpressionCache,
1709 ) -> Vec<BuiltinTableUpdate<&'static BuiltinTable>> {
1710 let mut builtin_table_updates = Vec::with_capacity(builtin_views.len());
1711 let (updates, additions): (Vec<_>, Vec<_>) =
1712 builtin_views
1713 .into_iter()
1714 .partition_map(|(view, item_id, gid)| {
1715 match retractions.system_object_mappings.remove(&item_id) {
1716 Some(entry) => Either::Left(entry),
1717 None => Either::Right((view, item_id, gid)),
1718 }
1719 });
1720
1721 for entry in updates {
1722 let item_id = entry.id();
1727 state.insert_entry(entry);
1728 builtin_table_updates.extend(state.pack_item_update(item_id, Diff::ONE));
1729 }
1730
1731 let mut handles = Vec::new();
1732 let mut awaiting_id_dependencies: BTreeMap<CatalogItemId, Vec<CatalogItemId>> =
1733 BTreeMap::new();
1734 let mut awaiting_name_dependencies: BTreeMap<String, Vec<CatalogItemId>> = BTreeMap::new();
1735 let mut awaiting_all = Vec::new();
1738 let mut completed_ids: BTreeSet<CatalogItemId> = BTreeSet::new();
1740 let mut completed_names: BTreeSet<String> = BTreeSet::new();
1741
1742 let mut views: BTreeMap<CatalogItemId, (&BuiltinView, GlobalId)> = additions
1744 .into_iter()
1745 .map(|(view, item_id, gid)| (item_id, (view, gid)))
1746 .collect();
1747 let item_ids: Vec<_> = views.keys().copied().collect();
1748
1749 let mut ready: VecDeque<CatalogItemId> = views.keys().cloned().collect();
1750 while !handles.is_empty() || !ready.is_empty() || !awaiting_all.is_empty() {
1751 if handles.is_empty() && ready.is_empty() {
1752 ready.extend(awaiting_all.drain(..));
1754 }
1755
1756 if !ready.is_empty() {
1758 let spawn_state = Arc::new(state.clone());
1759 while let Some(id) = ready.pop_front() {
1760 let (view, global_id) = views.get(&id).expect("must exist");
1761 let global_id = *global_id;
1762 let create_sql = view.create_sql();
1763 let versions = BTreeMap::new();
1765
1766 let span = info_span!(parent: None, "parse builtin view", name = view.name);
1767 OpenTelemetryContext::obtain().attach_as_parent_to(&span);
1768 let task_state = Arc::clone(&spawn_state);
1769 let cached_expr = local_expression_cache.remove_cached_expression(&global_id);
1770 let handle = mz_ore::task::spawn_blocking(
1771 || "parse view",
1772 move || {
1773 span.in_scope(|| {
1774 let res = task_state.parse_item_inner(
1775 global_id,
1776 &create_sql,
1777 &versions,
1778 None,
1779 false,
1780 None,
1781 cached_expr,
1782 None,
1783 );
1784 (id, global_id, res)
1785 })
1786 },
1787 );
1788 handles.push(handle);
1789 }
1790 }
1791
1792 let (selected, _idx, remaining) = future::select_all(handles).await;
1794 handles = remaining;
1795 let (id, global_id, res) = selected;
1796 let mut insert_cached_expr = |cached_expr| {
1797 if let Some(cached_expr) = cached_expr {
1798 local_expression_cache.insert_cached_expression(global_id, cached_expr);
1799 }
1800 };
1801 match res {
1802 Ok((item, uncached_expr)) => {
1803 if let Some((uncached_expr, optimizer_features)) = uncached_expr {
1804 local_expression_cache.insert_uncached_expression(
1805 global_id,
1806 uncached_expr,
1807 optimizer_features,
1808 RelationVersion::root(),
1809 );
1810 }
1811 let (view, _gid) = views.remove(&id).expect("must exist");
1813 let schema_id = state
1814 .ambient_schemas_by_name
1815 .get(view.schema)
1816 .unwrap_or_else(|| panic!("unknown ambient schema: {}", view.schema));
1817 let qname = QualifiedItemName {
1818 qualifiers: ItemQualifiers {
1819 database_spec: ResolvedDatabaseSpecifier::Ambient,
1820 schema_spec: SchemaSpecifier::Id(*schema_id),
1821 },
1822 item: view.name.into(),
1823 };
1824 let mut acl_items = vec![rbac::owner_privilege(
1825 mz_sql::catalog::ObjectType::View,
1826 MZ_SYSTEM_ROLE_ID,
1827 )];
1828 acl_items.extend_from_slice(&view.access);
1829
1830 state.insert_item(
1831 id,
1832 view.oid,
1833 qname,
1834 item,
1835 MZ_SYSTEM_ROLE_ID,
1836 PrivilegeMap::from_mz_acl_items(acl_items),
1837 );
1838
1839 let mut resolved_dependent_items = Vec::new();
1841 if let Some(dependent_items) = awaiting_id_dependencies.remove(&id) {
1842 resolved_dependent_items.extend(dependent_items);
1843 }
1844 let entry = state.get_entry(&id);
1845 let full_name = state.resolve_full_name(entry.name(), None).to_string();
1846 if let Some(dependent_items) = awaiting_name_dependencies.remove(&full_name) {
1847 resolved_dependent_items.extend(dependent_items);
1848 }
1849 ready.extend(resolved_dependent_items);
1850
1851 completed_ids.insert(id);
1852 completed_names.insert(full_name);
1853 }
1854 Err((
1856 AdapterError::PlanError(plan::PlanError::InvalidId(missing_dep)),
1857 cached_expr,
1858 )) => {
1859 insert_cached_expr(cached_expr);
1860 if completed_ids.contains(&missing_dep) {
1861 ready.push_back(id);
1862 } else {
1863 awaiting_id_dependencies
1864 .entry(missing_dep)
1865 .or_default()
1866 .push(id);
1867 }
1868 }
1869 Err((
1871 AdapterError::PlanError(plan::PlanError::Catalog(
1872 SqlCatalogError::UnknownItem(missing_dep),
1873 )),
1874 cached_expr,
1875 )) => {
1876 insert_cached_expr(cached_expr);
1877 match CatalogItemId::from_str(&missing_dep) {
1878 Ok(missing_dep) => {
1879 if completed_ids.contains(&missing_dep) {
1880 ready.push_back(id);
1881 } else {
1882 awaiting_id_dependencies
1883 .entry(missing_dep)
1884 .or_default()
1885 .push(id);
1886 }
1887 }
1888 Err(_) => {
1889 if completed_names.contains(&missing_dep) {
1890 ready.push_back(id);
1891 } else {
1892 awaiting_name_dependencies
1893 .entry(missing_dep)
1894 .or_default()
1895 .push(id);
1896 }
1897 }
1898 }
1899 }
1900 Err((
1901 AdapterError::PlanError(plan::PlanError::InvalidCast { .. }),
1902 cached_expr,
1903 )) => {
1904 insert_cached_expr(cached_expr);
1905 awaiting_all.push(id);
1906 }
1907 Err((e, _)) => {
1908 let (bad_view, _gid) = views.get(&id).expect("must exist");
1909 panic!(
1910 "internal error: failed to load bootstrap view:\n\
1911 {name}\n\
1912 error:\n\
1913 {e:?}\n\n\
1914 Make sure that the schema name is specified in the builtin view's create sql statement.
1915 ",
1916 name = bad_view.name,
1917 )
1918 }
1919 }
1920 }
1921
1922 assert!(awaiting_id_dependencies.is_empty());
1923 assert!(
1924 awaiting_name_dependencies.is_empty(),
1925 "awaiting_name_dependencies: {awaiting_name_dependencies:?}"
1926 );
1927 assert!(awaiting_all.is_empty());
1928 assert!(views.is_empty());
1929
1930 builtin_table_updates.extend(
1932 item_ids
1933 .into_iter()
1934 .flat_map(|id| state.pack_item_update(id, Diff::ONE)),
1935 );
1936
1937 builtin_table_updates
1938 }
1939
1940 fn insert_entry(&mut self, entry: CatalogEntry) {
1942 if !entry.id.is_system() {
1943 if let Some(cluster_id) = entry.item.cluster_id() {
1944 self.clusters_by_id
1945 .get_mut(&cluster_id)
1946 .expect("catalog out of sync")
1947 .bound_objects
1948 .insert(entry.id);
1949 };
1950 }
1951
1952 for u in entry.references().items() {
1953 match self.entry_by_id.get_mut(u) {
1954 Some(metadata) => metadata.referenced_by.push(entry.id()),
1955 None => panic!(
1956 "Catalog: missing dependent catalog item {} while installing {}",
1957 u,
1958 self.resolve_full_name(entry.name(), entry.conn_id())
1959 ),
1960 }
1961 }
1962 for u in entry.uses() {
1963 if u == entry.id() {
1966 continue;
1967 }
1968 match self.entry_by_id.get_mut(&u) {
1969 Some(metadata) => metadata.used_by.push(entry.id()),
1970 None => panic!(
1971 "Catalog: missing dependent catalog item {} while installing {}",
1972 u,
1973 self.resolve_full_name(entry.name(), entry.conn_id())
1974 ),
1975 }
1976 }
1977 for gid in entry.item.global_ids() {
1978 self.entry_by_global_id.insert(gid, entry.id());
1979 }
1980 let conn_id = entry.item().conn_id().unwrap_or(&SYSTEM_CONN_ID);
1981 if entry.name().qualifiers.schema_spec == SchemaSpecifier::Temporary {
1984 self.temporary_namespaces
1985 .ensure_schema(conn_id, entry.owner_id);
1986 }
1987 let schema = self.get_schema_mut(
1988 &entry.name().qualifiers.database_spec,
1989 &entry.name().qualifiers.schema_spec,
1990 conn_id,
1991 );
1992
1993 let prev_id = match entry.item() {
1994 CatalogItem::Func(_) => schema
1995 .functions
1996 .insert(entry.name().item.clone(), entry.id()),
1997 CatalogItem::Type(_) => schema.types.insert(entry.name().item.clone(), entry.id()),
1998 _ => schema.items.insert(entry.name().item.clone(), entry.id()),
1999 };
2000
2001 assert!(
2002 prev_id.is_none(),
2003 "builtin name collision on {:?}",
2004 entry.name().item.clone()
2005 );
2006
2007 self.entry_by_id.insert(entry.id(), entry.clone());
2008 }
2009
2010 fn insert_item(
2012 &mut self,
2013 id: CatalogItemId,
2014 oid: u32,
2015 name: QualifiedItemName,
2016 item: CatalogItem,
2017 owner_id: RoleId,
2018 privileges: PrivilegeMap,
2019 ) {
2020 let entry = CatalogEntry {
2021 item,
2022 name,
2023 id,
2024 oid,
2025 used_by: Vec::new(),
2026 referenced_by: Vec::new(),
2027 owner_id,
2028 privileges,
2029 };
2030
2031 self.insert_entry(entry);
2032 }
2033
2034 #[mz_ore::instrument(level = "trace")]
2035 fn drop_item(&mut self, id: CatalogItemId) -> CatalogEntry {
2036 let metadata = self.entry_by_id.remove(&id).expect("catalog out of sync");
2037 for u in metadata.references().items() {
2038 if let Some(dep_metadata) = self.entry_by_id.get_mut(u) {
2039 dep_metadata.referenced_by.retain(|u| *u != metadata.id())
2040 }
2041 }
2042 for u in metadata.uses() {
2043 if let Some(dep_metadata) = self.entry_by_id.get_mut(&u) {
2044 dep_metadata.used_by.retain(|u| *u != metadata.id())
2045 }
2046 }
2047 for gid in metadata.global_ids() {
2048 self.entry_by_global_id.remove(&gid);
2049 }
2050
2051 let conn_id = metadata.item().conn_id().unwrap_or(&SYSTEM_CONN_ID);
2052 let schema = self.get_schema_mut(
2053 &metadata.name().qualifiers.database_spec,
2054 &metadata.name().qualifiers.schema_spec,
2055 conn_id,
2056 );
2057 if metadata.item_type() == CatalogItemType::Type {
2058 schema
2059 .types
2060 .remove(&metadata.name().item)
2061 .expect("catalog out of sync");
2062 } else {
2063 assert_ne!(metadata.item_type(), CatalogItemType::Func);
2066
2067 schema
2068 .items
2069 .remove(&metadata.name().item)
2070 .expect("catalog out of sync");
2071 };
2072
2073 if !id.is_system() {
2074 if let Some(cluster_id) = metadata.item().cluster_id() {
2075 assert!(
2076 self.clusters_by_id
2077 .get_mut(&cluster_id)
2078 .expect("catalog out of sync")
2079 .bound_objects
2080 .remove(&id),
2081 "catalog out of sync"
2082 );
2083 }
2084 }
2085
2086 metadata
2087 }
2088
2089 fn insert_introspection_source_index(
2090 &mut self,
2091 cluster_id: ClusterId,
2092 log: &'static BuiltinLog,
2093 item_id: CatalogItemId,
2094 global_id: GlobalId,
2095 oid: u32,
2096 ) {
2097 let (index_name, index) =
2098 self.create_introspection_source_index(cluster_id, log, global_id);
2099 self.insert_item(
2100 item_id,
2101 oid,
2102 index_name,
2103 index,
2104 MZ_SYSTEM_ROLE_ID,
2105 PrivilegeMap::default(),
2106 );
2107 }
2108
2109 fn create_introspection_source_index(
2110 &self,
2111 cluster_id: ClusterId,
2112 log: &'static BuiltinLog,
2113 global_id: GlobalId,
2114 ) -> (QualifiedItemName, CatalogItem) {
2115 let source_name = FullItemName {
2116 database: RawDatabaseSpecifier::Ambient,
2117 schema: log.schema.into(),
2118 item: log.name.into(),
2119 };
2120 let index_name = format!("{}_{}_primary_idx", log.name, cluster_id);
2121 let mut index_name = QualifiedItemName {
2122 qualifiers: ItemQualifiers {
2123 database_spec: ResolvedDatabaseSpecifier::Ambient,
2124 schema_spec: SchemaSpecifier::Id(self.get_mz_introspection_schema_id()),
2125 },
2126 item: index_name.clone(),
2127 };
2128 index_name = self.find_available_name(index_name, &SYSTEM_CONN_ID);
2129 let index_item_name = index_name.item.clone();
2130 let (log_item_id, log_global_id) = self.resolve_builtin_log(log);
2131 let index = CatalogItem::Index(Index {
2132 global_id,
2133 on: log_global_id,
2134 keys: log
2135 .variant
2136 .index_by()
2137 .into_iter()
2138 .map(MirScalarExpr::column)
2139 .collect(),
2140 create_sql: index_sql(
2141 index_item_name,
2142 cluster_id,
2143 source_name,
2144 &log.variant.desc(),
2145 &log.variant.index_by(),
2146 ),
2147 conn_id: None,
2148 resolved_ids: [(log_item_id, log_global_id)].into_iter().collect(),
2149 cluster_id,
2150 is_retained_metrics_object: false,
2151 custom_logical_compaction_window: None,
2152 optimized_plan: None,
2153 physical_plan: None,
2154 dataflow_metainfo: None,
2155 });
2156 (index_name, index)
2157 }
2158
2159 fn insert_system_configuration(&mut self, name: &str, value: VarInput) -> Result<bool, Error> {
2164 Ok(Arc::make_mut(&mut self.system_configuration).set(name, value)?)
2165 }
2166
2167 fn remove_system_configuration(&mut self, name: &str) -> Result<bool, Error> {
2172 Ok(Arc::make_mut(&mut self.system_configuration).reset(name)?)
2173 }
2174}
2175
2176fn sort_updates(updates: Vec<StateUpdate>) -> Vec<StateUpdate> {
2184 fn push_update<T>(
2185 update: T,
2186 diff: StateDiff,
2187 retractions: &mut Vec<T>,
2188 additions: &mut Vec<T>,
2189 ) {
2190 match diff {
2191 StateDiff::Retraction => retractions.push(update),
2192 StateDiff::Addition => additions.push(update),
2193 }
2194 }
2195
2196 soft_assert_no_log!(
2197 updates.iter().map(|update| update.ts).all_equal(),
2198 "all timestamps should be equal: {updates:?}"
2199 );
2200 soft_assert_no_log!(
2201 {
2202 let mut dedup = BTreeSet::new();
2203 updates.iter().all(|update| dedup.insert(&update.kind))
2204 },
2205 "updates should be consolidated: {updates:?}"
2206 );
2207
2208 let mut pre_cluster_retractions = Vec::new();
2210 let mut pre_cluster_additions = Vec::new();
2211 let mut cluster_retractions = Vec::new();
2212 let mut cluster_additions = Vec::new();
2213 let mut builtin_item_updates = Vec::new();
2214 let mut item_retractions = Vec::new();
2215 let mut item_additions = Vec::new();
2216 let mut post_item_retractions = Vec::new();
2217 let mut post_item_additions = Vec::new();
2218 for update in updates {
2219 let diff = update.diff.clone();
2220 match update.kind {
2221 StateUpdateKind::Role(_)
2222 | StateUpdateKind::RoleAuth(_)
2223 | StateUpdateKind::Database(_)
2224 | StateUpdateKind::Schema(_)
2225 | StateUpdateKind::DefaultPrivilege(_)
2226 | StateUpdateKind::SystemPrivilege(_)
2227 | StateUpdateKind::SystemConfiguration(_)
2228 | StateUpdateKind::NetworkPolicy(_) => push_update(
2229 update,
2230 diff,
2231 &mut pre_cluster_retractions,
2232 &mut pre_cluster_additions,
2233 ),
2234 StateUpdateKind::Cluster(_)
2235 | StateUpdateKind::ClusterSystemConfiguration(_)
2236 | StateUpdateKind::IntrospectionSourceIndex(_)
2237 | StateUpdateKind::ClusterReplica(_)
2238 | StateUpdateKind::ReplicaSystemConfiguration(_) => push_update(
2239 update,
2240 diff,
2241 &mut cluster_retractions,
2242 &mut cluster_additions,
2243 ),
2244 StateUpdateKind::SystemObjectMapping(system_object_mapping) => {
2245 builtin_item_updates.push((system_object_mapping, update.ts, update.diff))
2246 }
2247 StateUpdateKind::Item(item) => push_update(
2248 (item, update.ts, update.diff),
2249 diff,
2250 &mut item_retractions,
2251 &mut item_additions,
2252 ),
2253 StateUpdateKind::Comment(_)
2254 | StateUpdateKind::SourceReferences(_)
2255 | StateUpdateKind::AuditLog(_)
2256 | StateUpdateKind::StorageCollectionMetadata(_)
2257 | StateUpdateKind::UnfinalizedShard(_) => push_update(
2258 update,
2259 diff,
2260 &mut post_item_retractions,
2261 &mut post_item_additions,
2262 ),
2263 }
2264 }
2265
2266 let builtin_item_updates = builtin_item_updates
2269 .into_iter()
2270 .map(|(system_object_mapping, ts, diff)| {
2271 let idx = BUILTIN_LOOKUP
2272 .get(&system_object_mapping.description)
2273 .expect("missing builtin")
2274 .0;
2275 (idx, system_object_mapping, ts, diff)
2276 })
2277 .sorted_by_key(|(idx, _, _, _)| *idx)
2278 .map(|(_, system_object_mapping, ts, diff)| (system_object_mapping, ts, diff));
2279
2280 let mut builtin_source_retractions = Vec::new();
2284 let mut builtin_source_additions = Vec::new();
2285 let mut other_builtin_retractions = Vec::new();
2286 let mut other_builtin_additions = Vec::new();
2287 for (builtin_item_update, ts, diff) in builtin_item_updates {
2288 let object_type = builtin_item_update.description.object_type;
2289 let update = StateUpdate {
2290 kind: StateUpdateKind::SystemObjectMapping(builtin_item_update),
2291 ts,
2292 diff,
2293 };
2294 if object_type == CatalogItemType::Source {
2295 push_update(
2296 update,
2297 diff,
2298 &mut builtin_source_retractions,
2299 &mut builtin_source_additions,
2300 );
2301 } else {
2302 push_update(
2303 update,
2304 diff,
2305 &mut other_builtin_retractions,
2306 &mut other_builtin_additions,
2307 );
2308 }
2309 }
2310
2311 fn sort_items_topological(items: &mut Vec<(mz_catalog::durable::Item, Timestamp, StateDiff)>) {
2317 tracing::debug!(?items, "sorting items by dependencies");
2318
2319 let key_fn = |item: &(mz_catalog::durable::Item, _, _)| item.0.id;
2320 let dependencies_fn = |item: &(mz_catalog::durable::Item, _, _)| {
2321 let statement = mz_sql::parse::parse(&item.0.create_sql)
2322 .expect("valid create_sql")
2323 .into_element()
2324 .ast;
2325 mz_sql::names::dependencies(&statement).expect("failed to find dependencies of item")
2326 };
2327 sort_topological(items, key_fn, dependencies_fn);
2328 }
2329
2330 fn sort_item_updates(
2346 item_updates: Vec<(mz_catalog::durable::Item, Timestamp, StateDiff)>,
2347 ) -> VecDeque<(mz_catalog::durable::Item, Timestamp, StateDiff)> {
2348 let mut types = Vec::new();
2351 let mut funcs = Vec::new();
2354 let mut secrets = Vec::new();
2355 let mut connections = Vec::new();
2356 let mut sources = Vec::new();
2357 let mut tables = Vec::new();
2358 let mut derived_items = Vec::new();
2359 let mut sinks = Vec::new();
2360 for update in item_updates {
2361 match update.0.item_type() {
2362 CatalogItemType::Type => types.push(update),
2363 CatalogItemType::Func => funcs.push(update),
2364 CatalogItemType::Secret => secrets.push(update),
2365 CatalogItemType::Connection => connections.push(update),
2366 CatalogItemType::Source => sources.push(update),
2367 CatalogItemType::Table => tables.push(update),
2368 CatalogItemType::View
2369 | CatalogItemType::MaterializedView
2370 | CatalogItemType::Index
2371 | CatalogItemType::MetricSink => derived_items.push(update),
2372 CatalogItemType::Sink => sinks.push(update),
2373 }
2374 }
2375
2376 sort_items_topological(&mut connections);
2380 sort_items_topological(&mut derived_items);
2381
2382 for group in [
2384 &mut types,
2385 &mut funcs,
2386 &mut secrets,
2387 &mut sources,
2388 &mut tables,
2389 &mut sinks,
2390 ] {
2391 group.sort_by_key(|(item, _, _)| item.id);
2392 }
2393
2394 iter::empty()
2395 .chain(types)
2396 .chain(funcs)
2397 .chain(secrets)
2398 .chain(connections)
2399 .chain(sources)
2400 .chain(tables)
2401 .chain(derived_items)
2402 .chain(sinks)
2403 .collect()
2404 }
2405
2406 fn into_state_updates(
2410 item_updates: VecDeque<(mz_catalog::durable::Item, Timestamp, StateDiff)>,
2411 ) -> Vec<StateUpdate> {
2412 item_updates
2413 .into_iter()
2414 .map(|(item, ts, diff)| StateUpdate {
2415 kind: StateUpdateKind::Item(item),
2416 ts,
2417 diff,
2418 })
2419 .collect()
2420 }
2421 let item_retractions = into_state_updates(sort_item_updates(item_retractions));
2422 let item_additions = into_state_updates(sort_item_updates(item_additions));
2423
2424 iter::empty()
2426 .chain(post_item_retractions.into_iter().rev())
2428 .chain(item_retractions.into_iter().rev())
2429 .chain(other_builtin_retractions.into_iter().rev())
2430 .chain(cluster_retractions.into_iter().rev())
2431 .chain(builtin_source_retractions.into_iter().rev())
2432 .chain(pre_cluster_retractions.into_iter().rev())
2433 .chain(pre_cluster_additions)
2434 .chain(builtin_source_additions)
2435 .chain(cluster_additions)
2436 .chain(other_builtin_additions)
2437 .chain(item_additions)
2438 .chain(post_item_additions)
2439 .collect()
2440}
2441
2442enum ApplyState {
2447 BuiltinViewAdditions(Vec<(&'static BuiltinView, CatalogItemId, GlobalId)>),
2449 Items(Vec<StateUpdate>),
2455 Updates(Vec<StateUpdate>),
2457}
2458
2459impl ApplyState {
2460 fn new(update: StateUpdate) -> Self {
2461 use StateUpdateKind::*;
2462 match &update.kind {
2463 SystemObjectMapping(som)
2464 if som.description.object_type == CatalogItemType::View
2465 && update.diff == StateDiff::Addition =>
2466 {
2467 let view_addition = lookup_builtin_view_addition(som.clone());
2468 Self::BuiltinViewAdditions(vec![view_addition])
2469 }
2470
2471 IntrospectionSourceIndex(_) | SystemObjectMapping(_) | Item(_) => {
2472 Self::Items(vec![update])
2473 }
2474
2475 Role(_)
2476 | RoleAuth(_)
2477 | Database(_)
2478 | Schema(_)
2479 | DefaultPrivilege(_)
2480 | SystemPrivilege(_)
2481 | SystemConfiguration(_)
2482 | ClusterSystemConfiguration(_)
2483 | ReplicaSystemConfiguration(_)
2484 | Cluster(_)
2485 | NetworkPolicy(_)
2486 | ClusterReplica(_)
2487 | SourceReferences(_)
2488 | Comment(_)
2489 | AuditLog(_)
2490 | StorageCollectionMetadata(_)
2491 | UnfinalizedShard(_) => Self::Updates(vec![update]),
2492 }
2493 }
2494
2495 async fn apply(
2501 self,
2502 state: &mut CatalogState,
2503 retractions: &mut InProgressRetractions,
2504 local_expression_cache: &mut LocalExpressionCache,
2505 ) -> (
2506 Vec<BuiltinTableUpdate<&'static BuiltinTable>>,
2507 Vec<ParsedStateUpdate>,
2508 ) {
2509 match self {
2510 Self::BuiltinViewAdditions(builtin_view_additions) => {
2511 let restore = Arc::clone(&state.system_configuration);
2512 Arc::make_mut(&mut state.system_configuration).enable_for_item_parsing();
2513 let builtin_table_updates = CatalogState::parse_builtin_views(
2514 state,
2515 builtin_view_additions,
2516 retractions,
2517 local_expression_cache,
2518 )
2519 .await;
2520 state.system_configuration = restore;
2521 (builtin_table_updates, Vec::new())
2522 }
2523 Self::Items(updates) => state.with_enable_for_item_parsing(|state| {
2524 state
2525 .apply_updates_inner(updates, retractions, local_expression_cache)
2526 .expect("corrupt catalog")
2527 }),
2528 Self::Updates(updates) => state
2529 .apply_updates_inner(updates, retractions, local_expression_cache)
2530 .expect("corrupt catalog"),
2531 }
2532 }
2533
2534 async fn step(
2535 self,
2536 next: Self,
2537 state: &mut CatalogState,
2538 retractions: &mut InProgressRetractions,
2539 local_expression_cache: &mut LocalExpressionCache,
2540 ) -> (
2541 Self,
2542 (
2543 Vec<BuiltinTableUpdate<&'static BuiltinTable>>,
2544 Vec<ParsedStateUpdate>,
2545 ),
2546 ) {
2547 match (self, next) {
2548 (
2549 Self::BuiltinViewAdditions(mut builtin_view_additions),
2550 Self::BuiltinViewAdditions(next_builtin_view_additions),
2551 ) => {
2552 builtin_view_additions.extend(next_builtin_view_additions);
2554 (
2555 Self::BuiltinViewAdditions(builtin_view_additions),
2556 (Vec::new(), Vec::new()),
2557 )
2558 }
2559 (Self::Items(mut updates), Self::Items(next_updates)) => {
2560 updates.extend(next_updates);
2562 (Self::Items(updates), (Vec::new(), Vec::new()))
2563 }
2564 (Self::Updates(mut updates), Self::Updates(next_updates)) => {
2565 updates.extend(next_updates);
2567 (Self::Updates(updates), (Vec::new(), Vec::new()))
2568 }
2569 (apply_state, next_apply_state) => {
2570 let updates = apply_state
2572 .apply(state, retractions, local_expression_cache)
2573 .await;
2574 (next_apply_state, updates)
2575 }
2576 }
2577 }
2578}
2579
2580trait MutableMap<K, V> {
2583 fn insert(&mut self, key: K, value: V) -> Option<V>;
2584 fn remove(&mut self, key: &K) -> Option<V>;
2585}
2586
2587impl<K: Ord, V> MutableMap<K, V> for BTreeMap<K, V> {
2588 fn insert(&mut self, key: K, value: V) -> Option<V> {
2589 BTreeMap::insert(self, key, value)
2590 }
2591 fn remove(&mut self, key: &K) -> Option<V> {
2592 BTreeMap::remove(self, key)
2593 }
2594}
2595
2596impl<K: Ord + Clone, V: Clone> MutableMap<K, V> for imbl::OrdMap<K, V> {
2597 fn insert(&mut self, key: K, value: V) -> Option<V> {
2598 imbl::OrdMap::insert(self, key, value)
2599 }
2600 fn remove(&mut self, key: &K) -> Option<V> {
2601 imbl::OrdMap::remove(self, key)
2602 }
2603}
2604
2605fn apply_inverted_lookup<K, V>(map: &mut impl MutableMap<K, V>, key: &K, value: V, diff: StateDiff)
2610where
2611 K: Ord + Clone + Debug,
2612 V: PartialEq + Debug,
2613{
2614 match diff {
2615 StateDiff::Retraction => {
2616 let prev = map.remove(key);
2617 assert_eq!(
2618 prev,
2619 Some(value),
2620 "retraction does not match existing value: {key:?}"
2621 );
2622 }
2623 StateDiff::Addition => {
2624 let prev = map.insert(key.clone(), value);
2625 assert_eq!(
2626 prev, None,
2627 "values must be explicitly retracted before inserting a new value: {key:?}"
2628 );
2629 }
2630 }
2631}
2632
2633fn apply_with_update<K, V, D>(
2636 map: &mut impl MutableMap<K, V>,
2637 durable: D,
2638 key_fn: impl FnOnce(&D) -> K,
2639 diff: StateDiff,
2640 retractions: &mut BTreeMap<D::Key, V>,
2641) where
2642 K: Ord,
2643 V: UpdateFrom<D> + PartialEq + Debug,
2644 D: DurableType,
2645 D::Key: Ord,
2646{
2647 match diff {
2648 StateDiff::Retraction => {
2649 let mem_key = key_fn(&durable);
2650 let value = map
2651 .remove(&mem_key)
2652 .expect("retraction does not match existing value: {key:?}");
2653 let durable_key = durable.into_key_value().0;
2654 retractions.insert(durable_key, value);
2655 }
2656 StateDiff::Addition => {
2657 let mem_key = key_fn(&durable);
2658 let durable_key = durable.key();
2659 let value = match retractions.remove(&durable_key) {
2660 Some(mut retraction) => {
2661 retraction.update_from(durable);
2662 retraction
2663 }
2664 None => durable.into(),
2665 };
2666 let prev = map.insert(mem_key, value);
2667 assert_eq!(
2668 prev, None,
2669 "values must be explicitly retracted before inserting a new value"
2670 );
2671 }
2672 }
2673}
2674
2675fn lookup_builtin_view_addition(
2677 mapping: SystemObjectMapping,
2678) -> (&'static BuiltinView, CatalogItemId, GlobalId) {
2679 let (_, builtin) = BUILTIN_LOOKUP
2680 .get(&mapping.description)
2681 .expect("missing builtin view");
2682 let Builtin::View(view) = builtin else {
2683 unreachable!("programming error, expected BuiltinView found {builtin:?}");
2684 };
2685
2686 (
2687 view,
2688 mapping.unique_identifier.catalog_id,
2689 mapping.unique_identifier.global_id,
2690 )
2691}