1use std::collections::{BTreeMap, BTreeSet};
32use std::sync::{Arc, LazyLock};
33
34use anyhow::bail;
35use futures::FutureExt;
36use futures::future::BoxFuture;
37use mz_build_info::{BuildInfo, DUMMY_BUILD_INFO};
38use mz_catalog::builtin::{
39 BUILTIN_LOOKUP, Builtin, Fingerprint, MZ_CATALOG_RAW, MZ_CATALOG_RAW_DESCRIPTION,
40 MZ_CLUSTER_REPLICA_FRONTIERS_DESCRIPTION, MZ_OBJECT_ARRANGEMENT_SIZE_HISTORY_DESCRIPTION,
41 MZ_OBJECT_HYDRATION_HISTORY, MZ_OBJECT_HYDRATION_HISTORY_DESCRIPTION,
42 MZ_REPLICA_HYDRATION_HISTORY, MZ_REPLICA_HYDRATION_HISTORY_DESCRIPTION,
43 MZ_STORAGE_USAGE_BY_SHARD, MZ_STORAGE_USAGE_BY_SHARD_DESCRIPTION,
44 RUNTIME_ALTERABLE_FINGERPRINT_SENTINEL,
45};
46use mz_catalog::config::BuiltinItemMigrationConfig;
47use mz_catalog::durable::objects::SystemObjectUniqueIdentifier;
48use mz_catalog::durable::{SystemObjectDescription, SystemObjectMapping, Transaction};
49use mz_catalog::memory::error::{Error, ErrorKind};
50use mz_ore::soft_assert_or_log;
51use mz_persist_client::cfg::USE_CRITICAL_SINCE_CATALOG;
52use mz_persist_client::critical::{Opaque, SinceHandle};
53use mz_persist_client::read::ReadHandle;
54use mz_persist_client::schema::CaESchema;
55use mz_persist_client::write::WriteHandle;
56use mz_persist_client::{Diagnostics, PersistClient};
57use mz_persist_types::ShardId;
58use mz_persist_types::codec_impls::{ShardIdSchema, UnitSchema};
59use mz_persist_types::schema::backward_compatible;
60use mz_repr::namespaces::{MZ_CATALOG_SCHEMA, MZ_INTERNAL_SCHEMA};
61use mz_repr::{CatalogItemId, GlobalId, Timestamp};
62use mz_sql::catalog::{CatalogItemType, NameReference};
63use mz_storage_client::controller::StorageTxn;
64use mz_storage_types::StorageDiff;
65use mz_storage_types::sources::SourceData;
66use semver::Version;
67use timely::progress::Antichain;
68use tracing::{debug, info};
69
70use crate::catalog::migrate::get_migration_version;
71
72static MIGRATIONS: LazyLock<Vec<MigrationStep>> = LazyLock::new(|| {
88 vec![
89 MigrationStep::replacement(
90 "0.149.0",
91 CatalogItemType::Source,
92 MZ_INTERNAL_SCHEMA,
93 "mz_sink_statistics_raw",
94 ),
95 MigrationStep::replacement(
96 "0.149.0",
97 CatalogItemType::Source,
98 MZ_INTERNAL_SCHEMA,
99 "mz_source_statistics_raw",
100 ),
101 MigrationStep::evolution(
102 "0.159.0",
103 CatalogItemType::Source,
104 MZ_INTERNAL_SCHEMA,
105 "mz_cluster_replica_metrics_history",
106 ),
107 MigrationStep::replacement(
108 "26.18.0-dev.0",
109 CatalogItemType::MaterializedView,
110 MZ_CATALOG_SCHEMA,
111 "mz_databases",
112 ),
113 MigrationStep::replacement(
114 "26.19.0-dev.0",
115 CatalogItemType::MaterializedView,
116 MZ_CATALOG_SCHEMA,
117 "mz_schemas",
118 ),
119 MigrationStep::replacement(
120 "26.19.0-dev.0",
121 CatalogItemType::MaterializedView,
122 MZ_CATALOG_SCHEMA,
123 "mz_role_members",
124 ),
125 MigrationStep::replacement(
126 "26.19.0-dev.0",
127 CatalogItemType::MaterializedView,
128 MZ_INTERNAL_SCHEMA,
129 "mz_network_policies",
130 ),
131 MigrationStep::replacement(
132 "26.19.0-dev.0",
133 CatalogItemType::MaterializedView,
134 MZ_INTERNAL_SCHEMA,
135 "mz_network_policy_rules",
136 ),
137 MigrationStep::replacement(
138 "26.19.0-dev.0",
139 CatalogItemType::MaterializedView,
140 MZ_INTERNAL_SCHEMA,
141 "mz_cluster_workload_classes",
142 ),
143 MigrationStep::replacement(
144 "26.19.0-dev.0",
145 CatalogItemType::MaterializedView,
146 MZ_INTERNAL_SCHEMA,
147 "mz_internal_cluster_replicas",
148 ),
149 MigrationStep::replacement(
150 "26.19.0-dev.0",
151 CatalogItemType::MaterializedView,
152 MZ_INTERNAL_SCHEMA,
153 "mz_pending_cluster_replicas",
154 ),
155 MigrationStep::replacement(
156 "26.20.0-dev.0",
157 CatalogItemType::MaterializedView,
158 MZ_CATALOG_SCHEMA,
159 "mz_materialized_views",
160 ),
161 MigrationStep::replacement(
162 "26.22.0-dev.0",
163 CatalogItemType::MaterializedView,
164 MZ_CATALOG_SCHEMA,
165 "mz_connections",
166 ),
167 MigrationStep::replacement(
168 "26.22.0-dev.0",
169 CatalogItemType::MaterializedView,
170 MZ_CATALOG_SCHEMA,
171 "mz_secrets",
172 ),
173 MigrationStep::replacement(
174 "26.27.0-dev.0",
175 CatalogItemType::MaterializedView,
176 MZ_CATALOG_SCHEMA,
177 "mz_sources",
178 ),
179 MigrationStep::replacement(
180 "26.29.0-dev.0",
181 CatalogItemType::MaterializedView,
182 MZ_CATALOG_SCHEMA,
183 "mz_indexes",
184 ),
185 MigrationStep::replacement(
186 "26.29.0-dev.0",
187 CatalogItemType::MaterializedView,
188 MZ_CATALOG_SCHEMA,
189 "mz_roles",
190 ),
191 MigrationStep::replacement(
192 "26.29.0-dev.0",
193 CatalogItemType::MaterializedView,
194 MZ_CATALOG_SCHEMA,
195 "mz_role_parameters",
196 ),
197 MigrationStep::replacement(
202 "26.30.0-dev.0",
203 CatalogItemType::MaterializedView,
204 MZ_CATALOG_SCHEMA,
205 "mz_indexes",
206 ),
207 MigrationStep::replacement(
208 "26.30.0-dev.0",
209 CatalogItemType::MaterializedView,
210 MZ_CATALOG_SCHEMA,
211 "mz_clusters",
212 ),
213 MigrationStep::replacement(
214 "26.30.0-dev.0",
215 CatalogItemType::MaterializedView,
216 MZ_CATALOG_SCHEMA,
217 "mz_cluster_replicas",
218 ),
219 MigrationStep::replacement(
220 "26.30.0-dev.0",
221 CatalogItemType::MaterializedView,
222 MZ_INTERNAL_SCHEMA,
223 "mz_cluster_schedules",
224 ),
225 MigrationStep::replacement(
226 "26.30.0-dev.0",
227 CatalogItemType::MaterializedView,
228 MZ_CATALOG_SCHEMA,
229 "mz_default_privileges",
230 ),
231 MigrationStep::replacement(
232 "26.30.0-dev.0",
233 CatalogItemType::MaterializedView,
234 MZ_CATALOG_SCHEMA,
235 "mz_system_privileges",
236 ),
237 MigrationStep::replacement(
241 "26.31.0-dev.0",
242 CatalogItemType::MaterializedView,
243 MZ_CATALOG_SCHEMA,
244 "mz_cluster_replicas",
245 ),
246 MigrationStep::replacement(
247 "26.32.0-dev.0",
248 CatalogItemType::MaterializedView,
249 MZ_INTERNAL_SCHEMA,
250 "mz_comments",
251 ),
252 MigrationStep::replacement(
257 "26.32.0-dev.0",
258 CatalogItemType::MaterializedView,
259 MZ_CATALOG_SCHEMA,
260 "mz_indexes",
261 ),
262 MigrationStep::replacement(
263 "26.33.0-dev.0",
264 CatalogItemType::MaterializedView,
265 MZ_CATALOG_SCHEMA,
266 "mz_audit_events",
267 ),
268 MigrationStep::replacement(
280 "26.34.0-dev.0",
281 CatalogItemType::MaterializedView,
282 MZ_CATALOG_SCHEMA,
283 "mz_indexes",
284 ),
285 MigrationStep::replacement(
290 "26.34.0-dev.0",
291 CatalogItemType::MaterializedView,
292 MZ_INTERNAL_SCHEMA,
293 "mz_postgres_sources",
294 ),
295 MigrationStep::replacement(
296 "26.34.0-dev.0",
297 CatalogItemType::MaterializedView,
298 MZ_CATALOG_SCHEMA,
299 "mz_kafka_sources",
300 ),
301 MigrationStep::replacement(
305 "26.37.0-dev.0",
306 CatalogItemType::MaterializedView,
307 MZ_INTERNAL_SCHEMA,
308 "mz_postgres_source_tables",
309 ),
310 MigrationStep::replacement(
311 "26.37.0-dev.0",
312 CatalogItemType::MaterializedView,
313 MZ_INTERNAL_SCHEMA,
314 "mz_mysql_source_tables",
315 ),
316 MigrationStep::replacement(
317 "26.37.0-dev.0",
318 CatalogItemType::MaterializedView,
319 MZ_INTERNAL_SCHEMA,
320 "mz_sql_server_source_tables",
321 ),
322 MigrationStep::replacement(
323 "26.37.0-dev.0",
324 CatalogItemType::MaterializedView,
325 MZ_INTERNAL_SCHEMA,
326 "mz_kafka_source_tables",
327 ),
328 MigrationStep::replacement(
332 "26.40.0-dev.0",
333 CatalogItemType::Table,
334 MZ_INTERNAL_SCHEMA,
335 "mz_type_pg_metadata",
336 ),
337 MigrationStep::replacement(
341 "26.37.0-dev.0",
342 CatalogItemType::MaterializedView,
343 MZ_CATALOG_SCHEMA,
344 "mz_kafka_connections",
345 ),
346 MigrationStep::replacement(
347 "26.37.0-dev.0",
348 CatalogItemType::MaterializedView,
349 MZ_CATALOG_SCHEMA,
350 "mz_ssh_tunnel_connections",
351 ),
352 MigrationStep::replacement(
353 "26.37.0-dev.0",
354 CatalogItemType::MaterializedView,
355 MZ_INTERNAL_SCHEMA,
356 "mz_aws_connections",
357 ),
358 MigrationStep::replacement(
359 "26.37.0-dev.0",
360 CatalogItemType::MaterializedView,
361 MZ_CATALOG_SCHEMA,
362 "mz_aws_privatelink_connections",
363 ),
364 MigrationStep::replacement(
370 "26.38.0-dev.0",
371 CatalogItemType::MaterializedView,
372 MZ_CATALOG_SCHEMA,
373 "mz_sinks",
374 ),
375 MigrationStep::replacement(
376 "26.38.0-dev.0",
377 CatalogItemType::MaterializedView,
378 MZ_CATALOG_SCHEMA,
379 "mz_kafka_sinks",
380 ),
381 MigrationStep::replacement(
382 "26.38.0-dev.0",
383 CatalogItemType::MaterializedView,
384 MZ_CATALOG_SCHEMA,
385 "mz_iceberg_sinks",
386 ),
387 MigrationStep::replacement(
390 "26.38.0-rc.2",
391 CatalogItemType::MaterializedView,
392 MZ_INTERNAL_SCHEMA,
393 "mz_cluster_reconfigurations",
394 ),
395 MigrationStep::replacement(
400 "26.39.0-dev.0",
401 CatalogItemType::MaterializedView,
402 MZ_CATALOG_SCHEMA,
403 "mz_audit_events",
404 ),
405 MigrationStep::replacement(
415 "26.40.0-dev.0",
416 CatalogItemType::MaterializedView,
417 MZ_CATALOG_SCHEMA,
418 "mz_indexes",
419 ),
420 MigrationStep::replacement(
426 "26.39.0-dev.0",
427 CatalogItemType::MaterializedView,
428 MZ_CATALOG_SCHEMA,
429 "mz_tables",
430 ),
431 MigrationStep::replacement(
432 "26.39.0-dev.0",
433 CatalogItemType::MaterializedView,
434 MZ_CATALOG_SCHEMA,
435 "mz_views",
436 ),
437 MigrationStep::replacement(
443 "26.40.0-dev.0",
444 CatalogItemType::MaterializedView,
445 MZ_CATALOG_SCHEMA,
446 "mz_indexes",
447 ),
448 MigrationStep::replacement(
449 "26.40.0-dev.0",
450 CatalogItemType::MaterializedView,
451 MZ_CATALOG_SCHEMA,
452 "mz_sources",
453 ),
454 MigrationStep::replacement(
460 "26.41.0-dev.0",
461 CatalogItemType::MaterializedView,
462 MZ_INTERNAL_SCHEMA,
463 "mz_object_dependencies",
464 ),
465 MigrationStep::evolution(
470 "26.43.0-dev.0",
471 CatalogItemType::Source,
472 MZ_INTERNAL_SCHEMA,
473 "mz_cluster_replica_metrics_history",
474 ),
475 MigrationStep::evolution(
476 "26.43.0-dev.0",
477 CatalogItemType::Table,
478 MZ_INTERNAL_SCHEMA,
479 "mz_replica_hydration_history",
480 ),
481 MigrationStep::replacement(
485 "26.44.0-dev.0",
486 CatalogItemType::MaterializedView,
487 MZ_INTERNAL_SCHEMA,
488 "mz_object_global_ids",
489 ),
490 MigrationStep::replacement(
494 "26.45.0-dev.0",
495 CatalogItemType::MaterializedView,
496 MZ_INTERNAL_SCHEMA,
497 "mz_source_references",
498 ),
499 ]
500});
501
502#[derive(Clone, Debug)]
504struct MigrationStep {
505 version: Version,
507 object: SystemObjectDescription,
509 mechanism: Mechanism,
511}
512
513impl MigrationStep {
514 fn evolution(version: &str, type_: CatalogItemType, schema: &str, name: &str) -> Self {
516 Self {
517 version: Version::parse(version).expect("valid"),
518 object: SystemObjectDescription {
519 schema_name: schema.into(),
520 object_type: type_,
521 object_name: name.into(),
522 },
523 mechanism: Mechanism::Evolution,
524 }
525 }
526
527 fn replacement(version: &str, type_: CatalogItemType, schema: &str, name: &str) -> Self {
529 Self {
530 version: Version::parse(version).expect("valid"),
531 object: SystemObjectDescription {
532 schema_name: schema.into(),
533 object_type: type_,
534 object_name: name.into(),
535 },
536 mechanism: Mechanism::Replacement,
537 }
538 }
539}
540
541#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
543#[allow(dead_code)]
544enum Mechanism {
545 Evolution,
550 Replacement,
554}
555
556pub(super) struct MigrationResult {
558 pub replaced_items: BTreeSet<CatalogItemId>,
560 pub cleanup_action: BoxFuture<'static, ()>,
562}
563
564impl Default for MigrationResult {
565 fn default() -> Self {
566 Self {
567 replaced_items: Default::default(),
568 cleanup_action: async {}.boxed(),
569 }
570 }
571}
572
573pub(super) async fn run(
579 build_info: &BuildInfo,
580 deploy_generation: u64,
581 txn: &mut Transaction<'_>,
582 config: BuiltinItemMigrationConfig,
583) -> Result<MigrationResult, Error> {
584 assert_eq!(config.read_only, txn.is_savepoint());
586
587 if *build_info == DUMMY_BUILD_INFO {
590 return Ok(MigrationResult::default());
591 }
592
593 let Some(durable_version) = get_migration_version(txn) else {
594 return Ok(MigrationResult::default());
596 };
597 let build_version = build_info.semver_version();
598
599 let collection_metadata = txn.get_collection_metadata();
600 let system_objects = txn
601 .get_system_object_mappings()
602 .map(|m| {
603 let object = m.description;
604 let global_id = m.unique_identifier.global_id;
605 let shard_id = collection_metadata.get(&global_id).copied();
606 let Some((_, builtin)) = BUILTIN_LOOKUP.get(&object) else {
607 panic!("missing builtin {object:?}");
608 };
609 let info = ObjectInfo {
610 global_id,
611 shard_id,
612 builtin,
613 fingerprint: m.unique_identifier.fingerprint,
614 };
615 (object, info)
616 })
617 .collect();
618
619 let migration_shard = txn.get_builtin_migration_shard().expect("must exist");
620
621 let migration = Migration {
622 source_version: durable_version.clone(),
623 target_version: build_version.clone(),
624 deploy_generation,
625 system_objects,
626 migration_shard,
627 config,
628 };
629
630 let result = migration.run(&MIGRATIONS).await.map_err(|e| {
631 Error::new(ErrorKind::FailedBuiltinSchemaMigration {
632 last_seen_version: durable_version.to_string(),
633 this_version: build_version.to_string(),
634 cause: e.to_string(),
635 })
636 })?;
637
638 result.apply(txn);
639
640 let replaced_items = txn
641 .get_system_object_mappings()
642 .map(|m| m.unique_identifier)
643 .filter(|ids| result.new_shards.contains_key(&ids.global_id))
644 .map(|ids| ids.catalog_id)
645 .collect();
646
647 Ok(MigrationResult {
648 replaced_items,
649 cleanup_action: result.cleanup_action,
650 })
651}
652
653struct MigrationRunResult {
655 new_shards: BTreeMap<GlobalId, ShardId>,
656 new_fingerprints: BTreeMap<SystemObjectDescription, String>,
657 shards_to_finalize: BTreeSet<ShardId>,
658 cleanup_action: BoxFuture<'static, ()>,
659}
660
661impl Default for MigrationRunResult {
662 fn default() -> Self {
663 Self {
664 new_shards: BTreeMap::new(),
665 new_fingerprints: BTreeMap::new(),
666 shards_to_finalize: BTreeSet::new(),
667 cleanup_action: async {}.boxed(),
668 }
669 }
670}
671
672impl MigrationRunResult {
673 fn apply(&self, txn: &mut Transaction<'_>) {
675 let replaced_ids = self.new_shards.keys().copied().collect();
677 let old_metadata = txn.delete_collection_metadata(replaced_ids);
678 txn.insert_collection_metadata(self.new_shards.clone())
679 .expect("inserting unique shards IDs after deleting existing entries");
680
681 let mut unfinalized_shards: BTreeSet<_> =
683 old_metadata.into_iter().map(|(_, sid)| sid).collect();
684 unfinalized_shards.extend(self.shards_to_finalize.iter().copied());
685 txn.insert_unfinalized_shards(unfinalized_shards)
686 .expect("cannot fail");
687
688 let mappings = txn
690 .get_system_object_mappings()
691 .filter_map(|m| {
692 let fingerprint = self.new_fingerprints.get(&m.description)?;
693 Some(SystemObjectMapping {
694 description: m.description,
695 unique_identifier: SystemObjectUniqueIdentifier {
696 catalog_id: m.unique_identifier.catalog_id,
697 global_id: m.unique_identifier.global_id,
698 fingerprint: fingerprint.clone(),
699 },
700 })
701 })
702 .collect();
703 txn.set_system_object_mappings(mappings)
704 .expect("filtered existing mappings remain unique");
705 }
706}
707
708#[derive(Clone, Debug)]
710struct ObjectInfo {
711 global_id: GlobalId,
712 shard_id: Option<ShardId>,
713 builtin: &'static Builtin<NameReference>,
714 fingerprint: String,
715}
716
717struct Migration {
719 source_version: Version,
724 target_version: Version,
728 deploy_generation: u64,
730 system_objects: BTreeMap<SystemObjectDescription, ObjectInfo>,
732 migration_shard: ShardId,
734 config: BuiltinItemMigrationConfig,
736}
737
738fn participates_in_forced_migration(
740 builtin: &Builtin<NameReference>,
741 mechanism: Mechanism,
742) -> bool {
743 use Builtin::*;
744 match builtin {
745 Table(table) => {
757 **table != *MZ_STORAGE_USAGE_BY_SHARD
758 && (mechanism != Mechanism::Replacement
759 || (**table != *MZ_OBJECT_HYDRATION_HISTORY
760 && **table != *MZ_REPLICA_HYDRATION_HISTORY))
761 }
762 MaterializedView(..) => true,
763 Source(source) => **source != *MZ_CATALOG_RAW,
764 Log(..) | View(..) | Type(..) | Func(..) | Index(..) | Connection(..) => false,
765 }
766}
767
768impl Migration {
769 async fn run(self, steps: &[MigrationStep]) -> anyhow::Result<MigrationRunResult> {
770 info!(
771 deploy_generation = %self.deploy_generation,
772 "running builtin schema migration: {} -> {}",
773 self.source_version, self.target_version
774 );
775
776 self.validate_migration_steps(steps);
777
778 let force_migration = if self.source_version != self.target_version
781 && self.source_version.pre.as_str().starts_with("dev")
782 && self.config.force_migration.is_none()
783 {
784 Some("evolution".to_string())
785 } else {
786 self.config.force_migration.clone()
787 };
788
789 let (force, plan) = match force_migration.as_deref() {
790 None => (false, self.plan_migration(steps)),
791 Some("evolution") => (true, self.plan_forced_migration(Mechanism::Evolution)),
792 Some("replacement") => (true, self.plan_forced_migration(Mechanism::Replacement)),
793 Some(other) => panic!("unknown force migration mechanism: {other}"),
794 };
795 let plan = self.drop_shardless(plan);
796
797 if self.source_version == self.target_version && !force {
798 info!("skipping migration: already at target version");
799 return Ok(MigrationRunResult::default());
800 } else if self.source_version > self.target_version {
801 bail!("downgrade not supported");
802 }
803
804 if !self.config.read_only {
807 self.upgrade_migration_shard_version().await;
808 }
809
810 info!("executing migration plan: {plan:?}");
811
812 self.migrate_evolve(&plan.evolve).await?;
813 let new_shards = self.migrate_replace(&plan.replace).await?;
814
815 let mut migrated_objects = BTreeSet::new();
816 migrated_objects.extend(plan.evolve);
817 migrated_objects.extend(plan.replace);
818
819 let new_fingerprints = self.update_fingerprints(&migrated_objects)?;
820
821 let (shards_to_finalize, cleanup_action) = self.cleanup().await?;
822
823 Ok(MigrationRunResult {
824 new_shards,
825 new_fingerprints,
826 shards_to_finalize,
827 cleanup_action,
828 })
829 }
830
831 fn validate_migration_steps(&self, steps: &[MigrationStep]) {
835 for step in steps {
836 assert!(
837 step.version <= self.target_version,
838 "migration step version greater than target version: {} > {}",
839 step.version,
840 self.target_version,
841 );
842
843 let object = &step.object;
844
845 assert_ne!(
853 &*MZ_STORAGE_USAGE_BY_SHARD_DESCRIPTION, object,
854 "mz_storage_usage_by_shard cannot be migrated or else the table will be truncated"
855 );
856
857 assert_ne!(
861 &*MZ_OBJECT_ARRANGEMENT_SIZE_HISTORY_DESCRIPTION, object,
862 "mz_object_arrangement_size_history cannot be migrated or else the table will be truncated"
863 );
864
865 if step.mechanism == Mechanism::Replacement {
876 assert_ne!(
877 &*MZ_OBJECT_HYDRATION_HISTORY_DESCRIPTION, object,
878 "replacing mz_object_hydration_history clears it, see the comment above"
879 );
880 assert_ne!(
881 &*MZ_REPLICA_HYDRATION_HISTORY_DESCRIPTION, object,
882 "replacing mz_replica_hydration_history clears it, see the comment above"
883 );
884 }
885
886 assert_ne!(
889 &*MZ_CATALOG_RAW_DESCRIPTION, object,
890 "mz_catalog_raw cannot be migrated"
891 );
892
893 assert_ne!(
898 &*MZ_CLUSTER_REPLICA_FRONTIERS_DESCRIPTION, object,
899 "mz_cluster_replica_frontiers cannot be migrated or else the 0dt caught-up gate loses its live-frontier reference"
900 );
901
902 let Some(object_info) = self.system_objects.get(object) else {
903 panic!("migration step for non-existent builtin: {object:?}");
904 };
905
906 let builtin = object_info.builtin;
907 use Builtin::*;
908 assert!(
909 matches!(builtin, Table(..) | Source(..) | MaterializedView(..)),
910 "schema migration not supported for builtin: {builtin:?}",
911 );
912 }
913 }
914
915 fn plan_migration(&self, steps: &[MigrationStep]) -> Plan {
917 let steps = steps.iter().filter(|s| s.version > self.source_version);
919
920 let mut by_object = BTreeMap::new();
924 for step in steps {
925 if let Some(entry) = by_object.get_mut(&step.object) {
926 *entry = match (step.mechanism, *entry) {
927 (Mechanism::Evolution, Mechanism::Evolution) => Mechanism::Evolution,
928 (Mechanism::Replacement, _) | (_, Mechanism::Replacement) => {
929 Mechanism::Replacement
930 }
931 };
932 } else {
933 by_object.insert(step.object.clone(), step.mechanism);
934 }
935 }
936
937 let mut plan = Plan::default();
938 for (object, mechanism) in by_object {
939 match mechanism {
940 Mechanism::Evolution => plan.evolve.push(object),
941 Mechanism::Replacement => plan.replace.push(object),
942 }
943 }
944
945 plan
946 }
947
948 fn plan_forced_migration(&self, mechanism: Mechanism) -> Plan {
950 let objects = self
951 .system_objects
952 .iter()
953 .filter(|(_, info)| participates_in_forced_migration(info.builtin, mechanism))
954 .map(|(object, _)| object.clone())
955 .collect();
956
957 let mut plan = Plan::default();
958 match mechanism {
959 Mechanism::Evolution => plan.evolve = objects,
960 Mechanism::Replacement => plan.replace = objects,
961 }
962
963 plan
964 }
965
966 fn drop_shardless(&self, mut plan: Plan) -> Plan {
976 let has_shard = |object: &SystemObjectDescription| {
977 self.system_objects
978 .get(object)
979 .is_none_or(|info| info.shard_id.is_some())
980 };
981 plan.evolve.retain(has_shard);
982 plan.replace.retain(has_shard);
983 plan
984 }
985
986 async fn upgrade_migration_shard_version(&self) {
988 let persist = &self.config.persist_client;
989 let diagnostics = Diagnostics {
990 shard_name: "builtin_migration".to_string(),
991 handle_purpose: format!("migration shard upgrade @ {}", self.target_version),
992 };
993
994 persist
995 .upgrade_version::<migration_shard::Key, ShardId, Timestamp, StorageDiff>(
996 self.migration_shard,
997 diagnostics,
998 )
999 .await
1000 .expect("valid usage");
1001 }
1002
1003 async fn migrate_evolve(&self, objects: &[SystemObjectDescription]) -> anyhow::Result<()> {
1005 for object in objects {
1006 self.migrate_evolve_one(object).await?;
1007 }
1008 Ok(())
1009 }
1010
1011 async fn migrate_evolve_one(&self, object: &SystemObjectDescription) -> anyhow::Result<()> {
1012 let persist = &self.config.persist_client;
1013
1014 let Some(object_info) = self.system_objects.get(object) else {
1015 bail!("missing builtin {object:?}");
1016 };
1017 let id = object_info.global_id;
1018
1019 let Some(shard_id) = object_info.shard_id else {
1020 if self.config.read_only {
1025 bail!("missing shard ID for builtin {object:?} ({id})");
1026 } else {
1027 return Ok(());
1028 }
1029 };
1030
1031 let target_desc = match object_info.builtin {
1032 Builtin::Table(table) => &table.desc,
1033 Builtin::Source(source) => &source.desc,
1034 Builtin::MaterializedView(mv) => &mv.desc,
1035 _ => bail!("not a storage collection: {object:?}"),
1036 };
1037
1038 let diagnostics = Diagnostics {
1039 shard_name: id.to_string(),
1040 handle_purpose: format!("builtin schema migration @ {}", self.target_version),
1041 };
1042 let source_schema = persist
1043 .latest_schema::<SourceData, (), Timestamp, StorageDiff>(shard_id, diagnostics.clone())
1044 .await
1045 .expect("valid usage");
1046
1047 info!(?object, %id, %shard_id, ?source_schema, ?target_desc, "migrating by evolution");
1048
1049 if self.config.read_only {
1050 if let Some((_, source_desc, _)) = &source_schema {
1053 let old = mz_persist_types::columnar::data_type::<SourceData>(source_desc)?;
1054 let new = mz_persist_types::columnar::data_type::<SourceData>(target_desc)?;
1055 if backward_compatible(&old, &new).is_none() {
1056 bail!(
1057 "incompatible schema evolution for {object:?}: \
1058 {source_desc:?} -> {target_desc:?}"
1059 );
1060 }
1061 }
1062
1063 return Ok(());
1064 }
1065
1066 let (mut schema_id, mut source_desc) = match source_schema {
1067 Some((schema_id, source_desc, _)) => (schema_id, source_desc),
1068 None => {
1069 debug!(%id, %shard_id, "no previous schema found; registering initial one");
1074 let schema_id = persist
1075 .register_schema::<SourceData, (), Timestamp, StorageDiff>(
1076 shard_id,
1077 target_desc,
1078 &UnitSchema,
1079 diagnostics.clone(),
1080 )
1081 .await
1082 .expect("valid usage");
1083 if schema_id.is_some() {
1084 return Ok(());
1085 }
1086
1087 debug!(%id, %shard_id, "schema registration failed; falling back to CaES");
1088 let (schema_id, source_desc, _) = persist
1089 .latest_schema::<SourceData, (), Timestamp, StorageDiff>(
1090 shard_id,
1091 diagnostics.clone(),
1092 )
1093 .await
1094 .expect("valid usage")
1095 .expect("known to exist");
1096
1097 (schema_id, source_desc)
1098 }
1099 };
1100
1101 loop {
1102 debug!(%id, %shard_id, %schema_id, ?source_desc, ?target_desc, "attempting CaES");
1107 let result = persist
1108 .compare_and_evolve_schema::<SourceData, (), Timestamp, StorageDiff>(
1109 shard_id,
1110 schema_id,
1111 target_desc,
1112 &UnitSchema,
1113 diagnostics.clone(),
1114 )
1115 .await
1116 .expect("valid usage");
1117
1118 match result {
1119 CaESchema::Ok(schema_id) => {
1120 debug!(%id, %shard_id, %schema_id, "schema evolved successfully");
1121 break;
1122 }
1123 CaESchema::Incompatible => bail!(
1124 "incompatible schema evolution for {object:?}: \
1125 {source_desc:?} -> {target_desc:?}"
1126 ),
1127 CaESchema::ExpectedMismatch {
1128 schema_id: new_id,
1129 key,
1130 val: UnitSchema,
1131 } => {
1132 schema_id = new_id;
1133 source_desc = key;
1134 }
1135 }
1136 }
1137
1138 Ok(())
1139 }
1140
1141 async fn migrate_replace(
1143 &self,
1144 objects: &[SystemObjectDescription],
1145 ) -> anyhow::Result<BTreeMap<GlobalId, ShardId>> {
1146 if objects.is_empty() {
1147 return Ok(Default::default());
1148 }
1149
1150 let diagnostics = Diagnostics {
1151 shard_name: "builtin_migration".to_string(),
1152 handle_purpose: format!("builtin schema migration @ {}", self.target_version),
1153 };
1154 let (mut persist_write, mut persist_read) =
1155 self.open_migration_shard(diagnostics.clone()).await;
1156
1157 let mut ids_to_replace = BTreeSet::new();
1158 for object in objects {
1159 if let Some(info) = self.system_objects.get(object) {
1160 ids_to_replace.insert(info.global_id);
1161 } else {
1162 bail!("missing id for builtin {object:?}");
1163 }
1164 }
1165
1166 info!(?objects, ?ids_to_replace, "migrating by replacement");
1167
1168 let replaced_shards = loop {
1171 if let Some(shards) = self
1172 .try_get_or_insert_replacement_shards(
1173 &ids_to_replace,
1174 &mut persist_write,
1175 &mut persist_read,
1176 )
1177 .await?
1178 {
1179 break shards;
1180 }
1181 };
1182
1183 Ok(replaced_shards)
1184 }
1185
1186 async fn try_get_or_insert_replacement_shards(
1196 &self,
1197 ids_to_replace: &BTreeSet<GlobalId>,
1198 persist_write: &mut WriteHandle<migration_shard::Key, ShardId, Timestamp, StorageDiff>,
1199 persist_read: &mut ReadHandle<migration_shard::Key, ShardId, Timestamp, StorageDiff>,
1200 ) -> anyhow::Result<Option<BTreeMap<GlobalId, ShardId>>> {
1201 let upper = persist_write.fetch_recent_upper().await;
1202 let write_ts = *upper.as_option().expect("migration shard not sealed");
1203
1204 let mut ids_to_replace = ids_to_replace.clone();
1205 let mut replaced_shards = BTreeMap::new();
1206
1207 if let Some(read_ts) = write_ts.step_back() {
1215 let pred = |key: &migration_shard::Key| {
1216 key.build_version == self.target_version
1217 && key.deploy_generation == Some(self.deploy_generation)
1218 };
1219 if let Some(entries) = read_migration_shard(persist_read, read_ts, pred).await {
1220 for (key, shard_id) in entries {
1221 let id = GlobalId::System(key.global_id);
1222 if ids_to_replace.remove(&id) {
1223 replaced_shards.insert(id, shard_id);
1224 }
1225 }
1226
1227 debug!(
1228 %read_ts, ?replaced_shards, ?ids_to_replace,
1229 "found existing entries in migration shard",
1230 );
1231 }
1232
1233 if ids_to_replace.is_empty() {
1234 return Ok(Some(replaced_shards));
1235 }
1236 }
1237
1238 let mut updates = Vec::new();
1242 for id in ids_to_replace {
1243 let shard_id = ShardId::new();
1244 replaced_shards.insert(id, shard_id);
1245
1246 let GlobalId::System(global_id) = id else {
1247 bail!("attempt to migrate a non-system collection: {id}");
1248 };
1249 let key = migration_shard::Key {
1250 global_id,
1251 build_version: self.target_version.clone(),
1252 deploy_generation: Some(self.deploy_generation),
1253 };
1254 updates.push(((key, shard_id), write_ts, 1));
1255 }
1256
1257 let upper = Antichain::from_elem(write_ts);
1258 let new_upper = Antichain::from_elem(write_ts.step_forward());
1259 debug!(%write_ts, "attempting insert into migration shard");
1260 let result = persist_write
1261 .compare_and_append(updates, upper, new_upper)
1262 .await
1263 .expect("valid usage");
1264
1265 match result {
1266 Ok(()) => {
1267 debug!(
1268 %write_ts, ?replaced_shards,
1269 "successfully inserted into migration shard"
1270 );
1271 Ok(Some(replaced_shards))
1272 }
1273 Err(_mismatch) => Ok(None),
1274 }
1275 }
1276
1277 async fn open_migration_shard(
1279 &self,
1280 diagnostics: Diagnostics,
1281 ) -> (
1282 WriteHandle<migration_shard::Key, ShardId, Timestamp, StorageDiff>,
1283 ReadHandle<migration_shard::Key, ShardId, Timestamp, StorageDiff>,
1284 ) {
1285 let persist = &self.config.persist_client;
1286
1287 persist
1288 .open(
1289 self.migration_shard,
1290 Arc::new(migration_shard::KeySchema),
1291 Arc::new(ShardIdSchema),
1292 diagnostics,
1293 USE_CRITICAL_SINCE_CATALOG.get(persist.dyncfgs()),
1294 )
1295 .await
1296 .expect("valid usage")
1297 }
1298
1299 async fn open_migration_shard_since(
1301 &self,
1302 diagnostics: Diagnostics,
1303 ) -> SinceHandle<migration_shard::Key, ShardId, Timestamp, StorageDiff> {
1304 self.config
1305 .persist_client
1306 .open_critical_since(
1307 self.migration_shard,
1308 PersistClient::CONTROLLER_CRITICAL_SINCE,
1311 Opaque::encode(&i64::MIN),
1312 diagnostics.clone(),
1313 )
1314 .await
1315 .expect("valid usage")
1316 }
1317
1318 fn update_fingerprints(
1323 &self,
1324 migrated_items: &BTreeSet<SystemObjectDescription>,
1325 ) -> anyhow::Result<BTreeMap<SystemObjectDescription, String>> {
1326 let mut new_fingerprints = BTreeMap::new();
1327 for (object, object_info) in &self.system_objects {
1328 let id = object_info.global_id;
1329 let builtin = object_info.builtin;
1330
1331 let fingerprint = builtin.fingerprint();
1332 if fingerprint == object_info.fingerprint {
1333 continue; }
1335
1336 let migrated = migrated_items.contains(object);
1338 let ephemeral = matches!(
1340 builtin,
1341 Builtin::Log(_) | Builtin::View(_) | Builtin::Index(_),
1342 );
1343
1344 if migrated || ephemeral {
1345 new_fingerprints.insert(object.clone(), fingerprint);
1346 } else if builtin.runtime_alterable() {
1347 assert_eq!(
1350 object_info.fingerprint, RUNTIME_ALTERABLE_FINGERPRINT_SENTINEL,
1351 "fingerprint mismatch for runtime-alterable builtin {object:?} ({id})",
1352 );
1353 } else {
1354 panic!(
1355 "fingerprint mismatch for builtin {builtin:?} ({id}): {} != {}",
1356 fingerprint, object_info.fingerprint,
1357 );
1358 }
1359 }
1360
1361 Ok(new_fingerprints)
1362 }
1363
1364 async fn cleanup(&self) -> anyhow::Result<(BTreeSet<ShardId>, BoxFuture<'static, ()>)> {
1381 let noop_action = async {}.boxed();
1382 let noop_result = (BTreeSet::new(), noop_action);
1383
1384 if self.config.read_only {
1385 return Ok(noop_result);
1386 }
1387
1388 let diagnostics = Diagnostics {
1389 shard_name: "builtin_migration".to_string(),
1390 handle_purpose: "builtin schema migration cleanup".into(),
1391 };
1392 let (mut persist_write, mut persist_read) =
1393 self.open_migration_shard(diagnostics.clone()).await;
1394 let mut persist_since = self.open_migration_shard_since(diagnostics.clone()).await;
1395
1396 let upper = persist_write.fetch_recent_upper().await.clone();
1397 let write_ts = *upper.as_option().expect("migration shard not sealed");
1398 let Some(read_ts) = write_ts.step_back() else {
1399 return Ok(noop_result);
1400 };
1401
1402 let pred = |key: &migration_shard::Key| key.build_version < self.target_version;
1404 let Some(stale_entries) = read_migration_shard(&mut persist_read, read_ts, pred).await
1405 else {
1406 return Ok(noop_result);
1407 };
1408
1409 debug!(
1410 ?stale_entries,
1411 "cleaning migration shard up to version {}", self.target_version,
1412 );
1413
1414 let current_shards: BTreeMap<_, _> = self
1415 .system_objects
1416 .values()
1417 .filter_map(|o| o.shard_id.map(|shard_id| (o.global_id, shard_id)))
1418 .collect();
1419
1420 let mut shards_to_finalize = BTreeSet::new();
1421 let mut retractions = Vec::new();
1422 for (key, shard_id) in stale_entries {
1423 let gid = GlobalId::System(key.global_id);
1427 if current_shards.get(&gid) != Some(&shard_id) {
1428 shards_to_finalize.insert(shard_id);
1429 }
1430
1431 retractions.push(((key, shard_id), write_ts, -1));
1432 }
1433
1434 let cleanup_action = async move {
1435 if !retractions.is_empty() {
1436 let new_upper = Antichain::from_elem(write_ts.step_forward());
1437 let result = persist_write
1438 .compare_and_append(retractions, upper, new_upper)
1439 .await
1440 .expect("valid usage");
1441 match result {
1442 Ok(()) => debug!("cleaned up migration shard"),
1443 Err(mismatch) => debug!(?mismatch, "migration shard cleanup failed"),
1444 }
1445 }
1446 }
1447 .boxed();
1448
1449 let o = persist_since.opaque().clone();
1451 let new_since = Antichain::from_elem(read_ts);
1452 let result = persist_since
1453 .maybe_compare_and_downgrade_since(&o, (&o, &new_since))
1454 .await;
1455 soft_assert_or_log!(result.is_none_or(|r| r.is_ok()), "opaque mismatch");
1456
1457 Ok((shards_to_finalize, cleanup_action))
1458 }
1459}
1460
1461async fn read_migration_shard<P>(
1467 persist_read: &mut ReadHandle<migration_shard::Key, ShardId, Timestamp, StorageDiff>,
1468 read_ts: Timestamp,
1469 predicate: P,
1470) -> Option<Vec<(migration_shard::Key, ShardId)>>
1471where
1472 P: Fn(&migration_shard::Key) -> bool,
1473{
1474 let as_of = Antichain::from_elem(read_ts);
1475 let updates = persist_read.snapshot_and_fetch(as_of).await.ok()?;
1476
1477 assert!(
1478 updates.iter().all(|(_, _, diff)| *diff == 1),
1479 "migration shard contains invalid diffs: {updates:?}",
1480 );
1481
1482 let entries: Vec<_> = updates
1483 .into_iter()
1484 .map(|(data, _, _)| data)
1485 .filter(move |(key, _)| predicate(key))
1486 .collect();
1487
1488 (!entries.is_empty()).then_some(entries)
1489}
1490
1491#[derive(Debug, Default)]
1493struct Plan {
1494 evolve: Vec<SystemObjectDescription>,
1496 replace: Vec<SystemObjectDescription>,
1498}
1499
1500mod migration_shard {
1502 use std::fmt;
1503 use std::str::FromStr;
1504
1505 use arrow::array::{StringArray, StringBuilder};
1506 use bytes::{BufMut, Bytes};
1507 use mz_persist_types::Codec;
1508 use mz_persist_types::codec_impls::{
1509 SimpleColumnarData, SimpleColumnarDecoder, SimpleColumnarEncoder,
1510 };
1511 use mz_persist_types::columnar::Schema;
1512 use mz_persist_types::stats::NoneStats;
1513 use semver::Version;
1514 use serde::{Deserialize, Serialize};
1515
1516 #[derive(Debug, Clone, Eq, Ord, PartialEq, PartialOrd, Serialize, Deserialize)]
1517 pub(super) struct Key {
1518 pub(super) global_id: u64,
1519 pub(super) build_version: Version,
1520 pub(super) deploy_generation: Option<u64>,
1524 }
1525
1526 impl fmt::Display for Key {
1527 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1528 if self.deploy_generation.is_some() {
1529 let s = serde_json::to_string(self).expect("JSON serializable");
1531 f.write_str(&s)
1532 } else {
1533 write!(f, "{}-{}", self.global_id, self.build_version)
1535 }
1536 }
1537 }
1538
1539 impl FromStr for Key {
1540 type Err = String;
1541
1542 fn from_str(s: &str) -> Result<Self, String> {
1543 if let Ok(key) = serde_json::from_str(s) {
1545 return Ok(key);
1546 };
1547
1548 let parts: Vec<_> = s.splitn(2, '-').collect();
1550 let &[global_id, build_version] = parts.as_slice() else {
1551 return Err(format!("invalid Key '{s}'"));
1552 };
1553 let global_id = global_id.parse::<u64>().map_err(|e| e.to_string())?;
1554 let build_version = build_version
1555 .parse::<Version>()
1556 .map_err(|e| e.to_string())?;
1557 Ok(Key {
1558 global_id,
1559 build_version,
1560 deploy_generation: None,
1561 })
1562 }
1563 }
1564
1565 impl Default for Key {
1566 fn default() -> Self {
1567 Self {
1568 global_id: Default::default(),
1569 build_version: Version::new(0, 0, 0),
1570 deploy_generation: Some(0),
1571 }
1572 }
1573 }
1574
1575 impl Codec for Key {
1576 type Schema = KeySchema;
1577 type Storage = ();
1578
1579 fn codec_name() -> String {
1580 "TableKey".into()
1581 }
1582
1583 fn encode<B: BufMut>(&self, buf: &mut B) {
1584 buf.put(self.to_string().as_bytes())
1585 }
1586
1587 fn decode<'a>(buf: &'a [u8], _schema: &KeySchema) -> Result<Self, String> {
1588 let s = str::from_utf8(buf).map_err(|e| e.to_string())?;
1589 s.parse()
1590 }
1591
1592 fn encode_schema(_schema: &KeySchema) -> Bytes {
1593 Bytes::new()
1594 }
1595
1596 fn decode_schema(buf: &Bytes) -> Self::Schema {
1597 assert_eq!(*buf, Bytes::new());
1598 KeySchema
1599 }
1600 }
1601
1602 impl SimpleColumnarData for Key {
1603 type ArrowBuilder = StringBuilder;
1604 type ArrowColumn = StringArray;
1605
1606 fn goodbytes(builder: &Self::ArrowBuilder) -> usize {
1607 builder.values_slice().len()
1608 }
1609
1610 fn push(&self, builder: &mut Self::ArrowBuilder) {
1611 builder.append_value(&self.to_string());
1612 }
1613
1614 fn push_null(builder: &mut Self::ArrowBuilder) {
1615 builder.append_null();
1616 }
1617
1618 fn read(&mut self, idx: usize, column: &Self::ArrowColumn) {
1619 *self = column.value(idx).parse().expect("valid Key");
1620 }
1621 }
1622
1623 #[derive(Debug, PartialEq)]
1624 pub(super) struct KeySchema;
1625
1626 impl Schema<Key> for KeySchema {
1627 type ArrowColumn = StringArray;
1628 type Statistics = NoneStats;
1629 type Decoder = SimpleColumnarDecoder<Key>;
1630 type Encoder = SimpleColumnarEncoder<Key>;
1631
1632 fn encoder(&self) -> anyhow::Result<SimpleColumnarEncoder<Key>> {
1633 Ok(SimpleColumnarEncoder::default())
1634 }
1635
1636 fn decoder(&self, col: StringArray) -> anyhow::Result<SimpleColumnarDecoder<Key>> {
1637 Ok(SimpleColumnarDecoder::new(col))
1638 }
1639 }
1640}
1641
1642#[cfg(test)]
1643#[path = "builtin_schema_migration_tests.rs"]
1644mod tests;